如何在异步处理场景下使用Kafka Listener的Acknowledgment机制?
异步使用Spring Kafka Acknowledgment的解决方案
问题背景
我们通过Spring Kafka的@KafkaListener向外部系统推送数据,外部系统采用带回调的异步通信渠道。为避免等待外部系统响应时阻塞下一批数据推送,希望在回调线程中根据处理结果调用Acknowledgment的ack()或nack()方法。目前ack()能正常调用,但nack()在回调线程中执行失败;尝试将nack()延迟到下一次Listener触发时执行,又会导致最后一批数据的nack()必须等待新数据到来才能执行,无法满足需求。
已尝试配置:
- AckMode设置为
MANUAL、MANUAL_IMMEDIATE,同时关闭enable.auto.commit - 开启/关闭
asyncAcks配置
核心问题:是否可以异步使用Acknowledgment功能?
解决方案1:使用DeferredAcknowledgment(推荐)
Spring Kafka 2.3及以上版本提供了DeferredAcknowledgment,专门适配异步确认场景。它会将ack()/nack()请求封装后提交给容器任务队列,由容器线程统一处理,避免了外部线程直接操作Consumer的安全问题,可在回调线程中安全调用nack()。
步骤1:配置容器工厂
确保容器工厂开启异步确认,并设置AckMode为MANUAL:
@Bean public ConcurrentKafkaListenerContainerFactory<String, CreateUser> batchFactory() { ConcurrentKafkaListenerContainerFactory<String, CreateUser> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); ContainerProperties containerProps = factory.getContainerProperties(); containerProps.setAckMode(ContainerProperties.AckMode.MANUAL); containerProps.setAsyncAcks(true); // 开启异步确认支持 return factory; }
步骤2:修改Listener方法
将参数中的Acknowledgment替换为DeferredAcknowledgment,即可在异步回调中正常调用nack():
@KafkaListener( id = CREATE_USERS_LISTENER_ID, topics = CREATE_USER_TOPIC, containerFactory = KafkaConstants.BATCH_FACTORY, autoStartup = "false" ) void simpleAsyncCase(ConsumerRecords<String, CreateUser> consumerRecords, DeferredAcknowledgment ack) { CompletableFuture.supplyAsync(() -> { int idx = 0; String last = ""; // 模拟调用外部系统并处理失败场景 Iterator<ConsumerRecord<String, CreateUser>> recordIterator = consumerRecords.iterator(); while (recordIterator.hasNext()) { ConsumerRecord<String, CreateUser> consumerRecord = recordIterator.next(); last = KafkaCoreListenerIT.toString(consumerRecord); if (random.nextBoolean()) { System.out.printf("XXXX Processing records %s at index %d%n", last, idx); latch.countDown(); } else { System.out.printf("XXXX ask to NACK record %s at index %d%n", last, idx); return idx; } idx++; } System.out.printf("XXXX ask to ACK to record %s at index %d%n", last, idx); return null; }).whenComplete((nackIdx, t) -> { try { if (t != null) { t.printStackTrace(); System.out.printf("XXXX NACK the full batch%n"); ack.nack(0, Duration.ZERO); // 异步线程可安全调用 } else if (nackIdx != null) { System.out.printf("XXXX NACK record at index %d%n", nackIdx); ack.nack(nackIdx, Duration.ZERO); // 异步线程可安全调用 } else { System.out.printf("XXXX Acknowledge"); ack.acknowledge(); } } catch (RuntimeException e) { e.printStackTrace(); throw e; } }); }
解决方案2:手动管理偏移量
如果无法使用DeferredAcknowledgment,可以手动管理偏移量,通过获取Consumer对象在回调线程中提交成功的偏移量,失败则不提交(触发重新消费)。这种方式需要自行处理偏移量的记录和提交逻辑,复杂度较高。
修改Listener方法
@KafkaListener( id = CREATE_USERS_LISTENER_ID, topics = CREATE_USER_TOPIC, containerFactory = KafkaConstants.BATCH_FACTORY, autoStartup = "false" ) void simpleAsyncCase(ConsumerRecords<String, CreateUser> consumerRecords, Consumer<String, CreateUser> consumer) { CompletableFuture.supplyAsync(() -> { int idx = 0; String last = ""; Map<TopicPartition, OffsetAndMetadata> successOffsets = new HashMap<>(); // 模拟调用外部系统并处理失败场景 Iterator<ConsumerRecord<String, CreateUser>> recordIterator = consumerRecords.iterator(); while (recordIterator.hasNext()) { ConsumerRecord<String, CreateUser> consumerRecord = recordIterator.next(); last = KafkaCoreListenerIT.toString(consumerRecord); if (random.nextBoolean()) { System.out.printf("XXXX Processing records %s at index %d%n", last, idx); latch.countDown(); // 记录成功处理的偏移量(下一个待消费的位置) successOffsets.put( new TopicPartition(consumerRecord.topic(), consumerRecord.partition()), new OffsetAndMetadata(consumerRecord.offset() + 1) ); } else { System.out.printf("XXXX ask to NACK record %s at index %d%n", last, idx); return successOffsets; // 返回已成功的偏移量,剩余数据将重新消费 } idx++; } System.out.printf("XXXX ask to ACK to record %s at index %d%n", last, idx); return successOffsets; }).whenComplete((successOffsets, t) -> { try { if (t != null) { t.printStackTrace(); System.out.printf("XXXX NACK the full batch%n"); // 不提交偏移量,等待容器重新拉取数据 } else if (!successOffsets.isEmpty()) { // 异步提交成功的偏移量 consumer.commitAsync(successOffsets, (committedOffsets, e) -> { if (e != null) { e.printStackTrace(); } }); } } catch (RuntimeException e) { e.printStackTrace(); throw e; } }); }
注意事项
- 需关闭
enable.auto.commit,并将AckMode设置为MANUAL或NONE - 手动提交偏移量需处理提交失败的重试逻辑,避免数据丢失或重复消费
内容的提问来源于stack exchange,提问作者uftapjech
相关产品推荐
相关产品推荐

