【发布时间】:2021-10-30 16:05:10
【问题描述】:
我已经设置了一个 spring-kafka 消费者。它使用来自主题的 avro 数据,映射值并写入 CSV 文件。一旦文件长度为 25000 条记录或每 5 分钟(以先到者为准),我手动提交偏移量。
由于补丁/发布,我们重新启动应用时出现问题。
我有这样的方法:
@PreDestroy
public void destroy() {
LOGGER.info("shutting down");
writeCsv(true);
acknowledgment.acknowledge(); // this normally commits the current offset
LOGGER.info("package commited: " + acknowledgment.toString());
LOGGER.info("shutting down completed");
}
所以我在那里添加了一些记录器,这就是日志的外观:
08:05:47 INFO KafkaMessageListenerContainer$ListenerConsumer - myManualConsumer: Consumer stopped
08:05:47 INFO CsvWriter - shutting down
08:05:47 INFO CsvWriter - created file: FEEDBACK1630476236079.csv
08:05:47 INFO CsvWriter - package commited: Acknowledgment for ConsumerRecord(topic = feedback-topic, partition = 1, leaderEpoch = 17, offset = 544, CreateTime = 1630415419703, serialized key size = -1, serialized value size = 156)
08:05:47 INFO CsvWriter - shutting down completed
由于消费者在调用确认()方法之前停止工作,因此永远不会提交偏移量。日志中没有错误,并且在应用程序重新启动后我们得到了重复。
- 有没有办法在消费者关闭之前调用方法?
还有一个问题:
我想像这样对消费者设置一个过滤器:
if(event.getValue().equals("GOOD") {
addCsvRecord(event)
} else {
acknowledgement.acknowledge() //to let it read next event
假设我得到了偏移量 100 - GOOD 事件来了,我将它添加到 csv 文件中,文件等待更多记录并且尚未提交偏移量。 接下来出现 BAD 事件,将其过滤掉并立即提交偏移量 101。 然后文件到达超时,即将关闭并调用
acknowlegdment.acknowledge()
在偏移 100 处。
- 那里可能发生什么?可以提交之前的偏移量吗?
【问题讨论】:
-
检查确认模式是否设置为
AbstractMessageListenerContainer.AckMode.MANUAL_IMMEDIATE -
设置为手动。将其设置为 MANUAL_IMMEDIATE 是否可以解决问题 #1? MANUAL - 处理完最后一次轮询的所有结果后,确认排队并在一次操作中提交偏移量。 MANUAL_IMMEDIATE - 只要在侦听器线程上执行 ack,就会立即提交偏移量(同步或异步)。似乎是一个解决方案!
-
是的,它可能会。试试看
-
它并没有真正帮助。我将在每条记录后提交偏移量来解决这个问题。不过感谢您的帮助! :)
-
@PreDestroy在上下文生命周期中为时已晚;到那时容器已经停止;看我的回答。
标签: java spring-boot apache-kafka spring-kafka offset