【发布时间】:2019-08-16 10:48:49
【问题描述】:
我已经使用 Kafka 作为源构建了一个 Storm 示例拓扑。这是一个我需要解决的问题。
每次我杀死一个拓扑并重新启动它时,拓扑都会从头开始处理。
假设 Topology X 中的消息 A 已被 Topology 处理,然后我终止了该 topology。
现在当我再次提交拓扑时,消息A仍然存在主题X。再次处理。
是否有解决方案,也许是某种偏移管理来处理这种情况。
【问题讨论】:
我已经使用 Kafka 作为源构建了一个 Storm 示例拓扑。这是一个我需要解决的问题。
每次我杀死一个拓扑并重新启动它时,拓扑都会从头开始处理。
假设 Topology X 中的消息 A 已被 Topology 处理,然后我终止了该 topology。
现在当我再次提交拓扑时,消息A仍然存在主题X。再次处理。
是否有解决方案,也许是某种偏移管理来处理这种情况。
【问题讨论】:
由于我面临类似的问题,请利用它并询问。我有这样的代码:
KafkaTridentSpoutConfig.Builder kafkaSpoutConfigBuilder = KafkaTridentSpoutConfig.builder(bootstrapServers, topic);
kafkaSpoutConfigBuilder.setProp(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, fetchSizeBytes);
kafkaSpoutConfigBuilder.setProp(ConsumerConfig.GROUP_ID_CONFIG, clientId);
kafkaSpoutConfigBuilder.setProp(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
kafkaSpoutConfigBuilder.setProp(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
return new KafkaTridentSpoutOpaque(kafkaSpoutConfigBuilder.build());
但是每次我重新启动 Storm Local Cluster 时,都会从头开始读取消息。如果我直接在 Kafka 中检查特定组的偏移量,则没有滞后。就像不读取来自 Kafka 的偏移量一样。 使用 Kafka 2.8、Storm 2.2.0。 Storm 0.9.X 没有这个问题。 有什么想法吗?
谢谢!
【讨论】:
您不应该将storm-kafka 用于新代码,它已被弃用,因为底层客户端 API 在 Kafka 中已弃用,并从 2.0.0 中删除。请改用storm-kafka-client。
对于storm-kafka-client,您想设置一个组 id 和第一个轮询偏移策略。
KafkaSpoutConfig.builder(bootstrapServers, "your-topic")
.setProp(ConsumerConfig.GROUP_ID_CONFIG, "kafkaSpoutTestGroup")
.setFirstPollOffsetStrategy(UNCOMMITTED_EARLIEST)
.build();
上面的内容将使您的 spout 在您第一次启动时从最早的偏移量开始,然后如果您重新启动它,它将从中断的地方开始。 Kafka 在 spout 重新启动时使用 group id 来识别它,因此它可以取回存储的偏移量检查点。其他偏移策略的行为会有所不同,您可以查看 javadoc 中的 FirstPollOffsetStrategy 枚举。
spout 将定期检查它的距离,配置中还有一个设置来控制它。检查点由配置中的setProcessingGuarantee 设置控制,并且可以设置为至少一次(仅检查点确认的偏移量)、最多一次(spout 发出消息之前的检查点)和“任何时候"(定期检查点,忽略确认)。
【讨论】:
确保在创建 spoutconfig 时,它有一个固定的 spout id,通过它可以在重启后识别自己。
来自 Storm 官方网站:
重要提示:重新部署拓扑时,请确保设置 对于 SpoutConfig.zkRoot 和 SpoutConfig.id 没有修改,否则 spout 将无法读取其先前的消费者状态 来自 ZooKeeper 的信息(即偏移量)——这可能导致 意外行为和/或数据丢失,具体取决于您的用例。
【讨论】: