【问题标题】:Apache Storm with Kafka offset management带有 Kafka 偏移管理的 Apache Storm
【发布时间】:2019-08-16 10:48:49
【问题描述】:

我已经使用 Kafka 作为源构建了一个 Storm 示例拓扑。这是一个我需要解决的问题。

每次我杀死一个拓扑并重新启动它时,拓扑都会从头开始处理。

假设 Topology X 中的消息 A 已被 Topology 处理,然后我终止了该 topology。

现在当我再次提交拓扑时,消息A仍然存在主题X。再次处理。

是否有解决方案,也许是某种偏移管理来处理这种情况。

【问题讨论】:

    标签: apache-kafka apache-storm


    【解决方案1】:

    由于我面临类似的问题,请利用它并询问。我有这样的代码:

            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 没有这个问题。 有什么想法吗?

    谢谢!

    【讨论】:

    • 这并不能真正回答问题。如果您有其他问题,可以点击 提问。要在此问题有新答案时收到通知,您可以follow this question。一旦你有足够的reputation,你也可以add a bounty 来引起对这个问题的更多关注。 - From Review
    【解决方案2】:

    您不应该将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 发出消息之前的检查点)和“任何时候"(定期检查点,忽略确认)。

    查看 Storm https://github.com/apache/storm/blob/dc56e32f3dcdd9396a827a85029d60ed97474786/examples/storm-kafka-client-examples/src/main/java/org/apache/storm/kafka/spout/KafkaSpoutTopologyMainNamedTopics.java#L93 中包含的示例拓扑之一。

    【讨论】:

      【解决方案3】:

      确保在创建 spoutconfig 时,它有一个固定的 spout id,通过它可以在重启后识别自己。

      来自 Storm 官方网站:

      重要提示:重新部署拓扑时,请确保设置 对于 SpoutConfig.zkRoot 和 SpoutConfig.id 没有修改,否则 spout 将无法读取其先前的消费者状态 来自 ZooKeeper 的信息(即偏移量)——这可能导致 意外行为和/或数据丢失,具体取决于您的用例。

      【讨论】:

      • 我假设偏移量只存储在确认的消息中?
      • 正如@Stig 所说“spout 将定期检查它有多远,配置中还有一个设置来控制这一点。检查点由配置中的 setProcessingGuarantee 设置控制,并且可以设置为至少一次(仅检查点确认的偏移量)、最多一次(在 spout 发出消息之前的检查点)和“任何时间”(定期检查点,忽略确认”
      猜你喜欢
      • 2019-06-25
      • 2017-07-17
      • 1970-01-01
      • 2018-09-22
      • 2021-01-15
      • 2018-09-24
      • 1970-01-01
      • 2020-07-20
      • 2017-07-09
      相关资源
      最近更新 更多