You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在异步处理场景下使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 22:05:27