【发布时间】: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