【发布时间】:2021-03-20 03:19:47
【问题描述】:
我正在尝试使用以下代码将 Kafka 集成到 Storm Toplogy,但不幸的是 KafkaSpout 没有使用来自 Kafka-topic 的消息。在 Storm UI-Core 中,Emitted count 永远保持为 0。
String bootStrapServer = "10.20.10.238:9092";
String topic = "test.topic";
KafkaSpoutConfig.Builder spoutConfigBuilder = KafkaSpoutConfig.builder(bootStrapServer,topic);
spoutConfigBuilder.setProp(ConsumerConfig.RECEIVE_BUFFER_CONFIG,100*1024*1024);
spoutConfigBuilder.setProp(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG,100*1024*1024);
spoutConfigBuilder.setProcessingGuarantee(KafkaSpoutConfig.ProcessingGuarantee.AT_LEAST_ONCE);
Boolean readFromStart = true;
if(readFromStart) {
spoutConfigBuilder.setFirstPollOffsetStrategy(FirstPollOffsetStrategy.EARLIEST);
}
else {
spoutConfigBuilder.setFirstPollOffsetStrategy(FirstPollOffsetStrategy.LATEST);
}
KafkaSpout spout = new KafkaSpout(spoutConfigBuilder.build());
builder.setSpout("kafkaSpout", spout, 1);
// And a Bolt to see messages
builder.setBolt("fcBolt", new FcBolt(), 1).setNumTasks(1).shuffleGrouping("kafkaSpout");
但是当我尝试从 CLI 查看生成的消息时,我可以使用以下命令查看有关主题的所有消息:
bin/kafka-console-consumer.sh --topic test.topic --from-beginning --bootstrap-server 10.20.10.238:9092
Picked up _JAVA_OPTIONS: -Xmx128000m
test
test
test1
....
版本:
Storm : 2.2.0
Kafka : 2.13_2.6.0
在旧版本中,它工作正常!我错过了在较新版本中阅读的内容。
任何帮助表示赞赏。提前致谢!
【问题讨论】:
标签: java apache-kafka kafka-consumer-api apache-storm