storm消费kafka数据

2024-05-31 17:38
文章标签 数据 kafka 消费 storm

本文主要是介绍storm消费kafka数据,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

http://blog.csdn.net/tonylee0329/article/details/43016385
使用storm-kafka模块读取kafka中的数据,按照以下两步进行构建(我使用的版本是0.9.3)
1. 使用BrokerHosts接口来配置kafka broker host与partition的mapping信息;
2. 使用KafkaConfig来配置一些与kafka自身相关的选项,如fetchSizeBytes、socketTimeoutMs
下面分别介绍这两块的实现:

对于配置1,目前支持两种实现方式:zk配置、静态ip端口方式

第一种方式:Zk读取(比较常见)
[html] view plain copy
ZkHosts支持两种创建方式,
public ZkHosts(String brokerZkStr, String brokerZkPath)
//使用默认brokerZkPath:”/brokers”
public ZkHosts(String brokerZkStr)

通过这种方式访问的时候,经过60s会刷新一下host->partition的mapping

第二步:构建KafkaConfig对象
目前提供两种构造函数,
[html] view plain copy
public KafkaConfig(BrokerHosts hosts, String topic)
//clientId如果不想每次随机生成的话,就自己设置一个
public KafkaConfig(BrokerHosts hosts, String topic, String clientId)

代码参考:
[html] view plain copy
//这个地方其实就是kafka配置文件里边的zookeeper.connect这个参数,可以去那里拿过来。
String brokerZkStr = “10.100.90.201:2181/kafka_online_sample”;
String brokerZkPath = “/brokers”;
ZkHosts zkHosts = new ZkHosts(brokerZkStr, brokerZkPath);

    String topic = "mars-wap";  //以下:将offset汇报到哪个zk集群,相应配置  

// String offsetZkServers = “10.199.203.169”;
String offsetZkServers = “10.100.90.201”;
String offsetZkPort = “2181”;
List zkServersList = new ArrayList();
zkServersList.add(offsetZkServers);
//汇报offset信息的root路径
String offsetZkRoot = “/stormExample”;
//存储该spout id的消费offset信息,譬如以topoName来命名
String offsetZkId = “storm-example”;

    SpoutConfig kafkaConfig = new SpoutConfig(zkHosts, topic, offsetZkRoot, offsetZkId);  kafkaConfig.zkRoot = offsetZkRoot;  kafkaConfig.zkPort = Integer.parseInt(offsetZkPort);  kafkaConfig.zkServers = zkServersList;  kafkaConfig.id = offsetZkId;  kafkaConfig.scheme = new SchemeAsMultiScheme(new StringScheme());  KafkaSpout spout = new KafkaSpout(kafkaConfig);  TopologyBuilder builder = new TopologyBuilder();  builder.setSpout("spout", spout, 1);  builder.setBolt("bolt", new Bolt(), 1).shuffleGrouping("spout");  Config config = new Config();  LocalCluster cluster = new LocalCluster();  cluster.submitTopology("test", config, builder.createTopology());  // cluster submit.  

// try {
// StormSubmitter.submitTopology(“storm-kafka-example”,config,builder.createTopology());
// } catch (AlreadyAliveException e) {
// e.printStackTrace();
// } catch (InvalidTopologyException e) {
// e.printStackTrace();
// }

第二种方式:静态ip端口方式
[html] view plain copy
String kafkaHost = “10.100.90.201”;
Broker brokerForPartition0 = new Broker(kafkaHost);//localhost:9092
Broker brokerForPartition1 = new Broker(kafkaHost, 9092);//localhost:9092 but we specified the port explicitly
GlobalPartitionInformation partitionInfo = new GlobalPartitionInformation();
partitionInfo.addPartition(0, brokerForPartition0);//mapping form partition 0 to brokerForPartition0
partitionInfo.addPartition(1, brokerForPartition1);//mapping form partition 1 to brokerForPartition1
StaticHosts hosts = new StaticHosts(partitionInfo);

    String topic="mars-wap";  String offsetZkRoot ="/stormExample";  String offsetZkId="staticHost";  String offsetZkServers = "10.100.90.201";  String offsetZkPort = "2181";  List<String> zkServersList = new ArrayList<String>();  zkServersList.add(offsetZkServers);  SpoutConfig kafkaConfig = new SpoutConfig(hosts,topic,offsetZkRoot,offsetZkId);  kafkaConfig.zkPort = Integer.parseInt(offsetZkPort);  kafkaConfig.zkServers = zkServersList;  kafkaConfig.scheme = new SchemeAsMultiScheme(new StringScheme());  KafkaSpout spout = new KafkaSpout(kafkaConfig);  TopologyBuilder builder = new TopologyBuilder();  builder.setSpout("spout", spout, 1);  builder.setBolt("bolt", new Bolt(), 1).shuffleGrouping("spout");  Config config = new Config();  LocalCluster cluster = new LocalCluster();  cluster.submitTopology("test", config, builder.createTopology());  

完整的使用例子,见github源码
https://github.com/tonylee0329/storm-example/blob/master/src/main/java/org/tony/storm_kafka/common/

参考:
https://github.com/apache/storm/blob/v0.9.3/external/storm-kafka/README.md

https://github.com/tonylee0329/storm-example/blob/master/src/main/java/org/tony/storm_kafka/common/ZkTopology.java

Kafka之Consumer获取消费数据全过程图解
字数198 阅读557 评论0 喜欢1
这篇文章是作为:跟我学Kafka源码之Consumer分析 的补充材料,看过我们之前源码分析的同学可能知道。
本文将从客户端程序如何调用Consumer获取到最终Kafka消息的全过程以图解的方式作一个源码级别的梳理。

废话不多说,请图看

时序图

Business Process Model.jpg
流程图

20140809174809543.png
文章短小的目的是便于大家快速找到内容的核心加以理解,避免文章又臭又长抓不住重点。
对于Kafka技术,如果大家对此有任何疑问,请给我留言,我们可以深入探讨。

清晰的UML时序图在这里可以看:
http://dl2.iteye.com/upload/attachment/0115/5649/70a096f4-c649-3efd-84bb-2379927dee36.jpg

这篇关于storm消费kafka数据的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1018603

相关文章

Linux下利用select实现串口数据读取过程

《Linux下利用select实现串口数据读取过程》文章介绍Linux中使用select、poll或epoll实现串口数据读取,通过I/O多路复用机制在数据到达时触发读取,避免持续轮询,示例代码展示设... 目录示例代码(使用select实现)代码解释总结在 linux 系统里,我们可以借助 select、

C#使用iText获取PDF的trailer数据的代码示例

《C#使用iText获取PDF的trailer数据的代码示例》开发程序debug的时候,看到了PDF有个trailer数据,挺有意思,于是考虑用代码把它读出来,那么就用到我们常用的iText框架了,所... 目录引言iText 核心概念C# 代码示例步骤 1: 确保已安装 iText步骤 2: C# 代码程

Pandas处理缺失数据的方式汇总

《Pandas处理缺失数据的方式汇总》许多教程中的数据与现实世界中的数据有很大不同,现实世界中的数据很少是干净且同质的,本文我们将讨论处理缺失数据的一些常规注意事项,了解Pandas如何表示缺失数据,... 目录缺失数据约定的权衡Pandas 中的缺失数据None 作为哨兵值NaN:缺失的数值数据Panda

C++中处理文本数据char与string的终极对比指南

《C++中处理文本数据char与string的终极对比指南》在C++编程中char和string是两种用于处理字符数据的类型,但它们在使用方式和功能上有显著的不同,:本文主要介绍C++中处理文本数... 目录1. 基本定义与本质2. 内存管理3. 操作与功能4. 性能特点5. 使用场景6. 相互转换核心区别

python库pydantic数据验证和设置管理库的用途

《python库pydantic数据验证和设置管理库的用途》pydantic是一个用于数据验证和设置管理的Python库,它主要利用Python类型注解来定义数据模型的结构和验证规则,本文给大家介绍p... 目录主要特点和用途:Field数值验证参数总结pydantic 是一个让你能够 confidentl

JAVA实现亿级千万级数据顺序导出的示例代码

《JAVA实现亿级千万级数据顺序导出的示例代码》本文主要介绍了JAVA实现亿级千万级数据顺序导出的示例代码,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面... 前提:主要考虑控制内存占用空间,避免出现同时导出,导致主程序OOM问题。实现思路:A.启用线程池

SpringBoot分段处理List集合多线程批量插入数据方式

《SpringBoot分段处理List集合多线程批量插入数据方式》文章介绍如何处理大数据量List批量插入数据库的优化方案:通过拆分List并分配独立线程处理,结合Spring线程池与异步方法提升效率... 目录项目场景解决方案1.实体类2.Mapper3.spring容器注入线程池bejsan对象4.创建

PHP轻松处理千万行数据的方法详解

《PHP轻松处理千万行数据的方法详解》说到处理大数据集,PHP通常不是第一个想到的语言,但如果你曾经需要处理数百万行数据而不让服务器崩溃或内存耗尽,你就会知道PHP用对了工具有多强大,下面小编就... 目录问题的本质php 中的数据流处理:为什么必不可少生成器:内存高效的迭代方式流量控制:避免系统过载一次性

C#实现千万数据秒级导入的代码

《C#实现千万数据秒级导入的代码》在实际开发中excel导入很常见,现代社会中很容易遇到大数据处理业务,所以本文我就给大家分享一下千万数据秒级导入怎么实现,文中有详细的代码示例供大家参考,需要的朋友可... 目录前言一、数据存储二、处理逻辑优化前代码处理逻辑优化后的代码总结前言在实际开发中excel导入很

MyBatis-plus处理存储json数据过程

《MyBatis-plus处理存储json数据过程》文章介绍MyBatis-Plus3.4.21处理对象与集合的差异:对象可用内置Handler配合autoResultMap,集合需自定义处理器继承F... 目录1、如果是对象2、如果需要转换的是List集合总结对象和集合分两种情况处理,目前我用的MP的版本