【问题标题】:How to Partial commitSync when consume batch in kafka在kafka中消费批处理时如何部分commitSync
【发布时间】:2018-11-27 01:45:27
【问题描述】:

我们正在使用 kafka 进行批量消费。我们消费 X 消息将它们放在 MYSQL 上然后提交它们。

有时我们会部分插入 MYSQL(重复记录、其他故障等)

使用文档中的这个例子:

List<ConsumerRecord<String, String>> buffer = new ArrayList<>();
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(100);
        for (ConsumerRecord<String, String> record : records) {
            buffer.add(record);
        }
        if (buffer.size() >= minBatchSize) {
            insertIntoDb(buffer);
            consumer.commitSync();
            buffer.clear();
        }

我们希望只 commitSync 成功的记录,同时让 kafka 重播失败。

但我无法理解如何执行此操作,因为 api 在整个批次中只有 commitSync()。

想法?

【问题讨论】:

  • 您可以使用其他风格的commitSync(请参阅kafka.apache.org/21/javadoc/org/apache/kafka/clients/consumer/…)指定偏移量或将您的投票配置为仅获取1条消息。
  • 为什么不使用 Kafka Connect JDBC Sink 为您处理这个逻辑?
  • @cricket_007 kafka connect jdbc sink 是否在出现任何问题时重新交付?
  • 取决于问题,但在一般情况下,是的。并且它比依赖缓冲区大小更密切地跟踪偏移

标签: java apache-kafka kafka-consumer-api


【解决方案1】:

在 Kafka 中,您不提交特定记录,即您不能将偏移量 N 标记为已处理,将偏移量 N-1 标记为未处理。相反,通过提交偏移量 N,您表明您已经处理了最多 N 的所有记录。

在处理偏移 N 失败时可以做的事情:

  • 提交 N-1(使用 commitSync(java.util.Map&lt;TopicPartition,OffsetAndMetadata&gt; offsets))并重试处理偏移量 N,因为它仍在内存中。只有在 N 成功处理后,您才提交 N 并移动到更新的记录。

  • 假设您在 Kafka Connect 的 Sink 连接器中运行,在处理 N 失败时,您可以将记录转发到 Connect 的 Deal Letter Queue。否则将其推回另一个主题以供以后处理。这会暂时跳过偏移量 N(如果可以的话,您也可以放弃它)。

您也可以混合使用这两种方法,重试几次,但如果无法处理此记录,请保存/删除它并继续处理更新的记录。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-12-15
    • 1970-01-01
    • 2021-05-06
    • 2021-02-01
    • 2020-12-18
    相关资源
    最近更新 更多