【发布时间】:2020-09-01 09:52:25
【问题描述】:
我正在考虑在线程池中使用 Kafka Consumer。我提出了这种方法。现在它似乎工作正常,但我正在考虑缺点以及这种方法会带来什么问题。基本上我需要的是将记录处理与消费分离。此外,我需要有一个强有力的保证,即只有在处理完所有记录后才会提交。有人可以就如何更好地做到这一点提出建议或建议吗?
final var consumer = new KafkaConsumer<String, String>(props);
consumer.subscribe(topics);
final var threadPool = Executors.newFixedThreadPool(32);
while(true) {
ConsumerRecords<String, String> records;
synchronized (consumer) {
records = consumer.poll(Duration.ofMillis(100));
}
CompletableFuture.runAsync(this::processTask, threadPool).thenRun(() -> {
synchronized (consumer) {
consumer.commitSync();
}
});
}
【问题讨论】:
-
这里可以写完整的代码吗?我假设会有一个 while 循环轮询 kafka 消费者以连续获取您希望以异步方式处理的记录。
-
是的,我更新了问题。
标签: java apache-kafka kafka-consumer-api