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

Kafka producer调用flush()返回时是否所有发送回调都已执行?

问题解答

核心结论

是的,当KafkaProducer.flush()方法返回时,可以100%保证所有之前通过send()提交的记录的回调都已经被触发执行完毕。

原理说明

KafkaProducer的flush()的设计语义就是阻塞当前线程,直到缓冲区中所有待发送的记录全部完成发送流程:不管是发送成功收到集群确认,还是发送失败抛出异常,对应的Callback.onCompletion()方法都会执行完成后,flush()才会结束阻塞返回。
你已经配置了acks=all参数,此时回调触发时如果没有异常,就代表对应记录已经被集群所有同步副本成功持久化存储,完全符合你校验发送状态的需求。

现有代码的注意点

你的实现目前存在两个可以优化的问题:

  • 自定义回调类的拼写错误:implememts需要修正为implements,否则会编译失败
  • 单回调实例共用的缺陷:你所有的send()调用都复用了同一个MyCallback实例,内部仅用单个AtomicReference存储异常,当有多条记录发送失败时,只有最后一个异常会被保留,之前的异常会被覆盖,无法感知全部失败情况。如果需要完整捕获所有发送失败的记录,建议改用线程安全的列表存储所有异常实例,check()方法时判断列表是否非空再抛出聚合异常即可。

改进后回调示例(可选)

public class MyCallback implements Callback {
    private final CopyOnWriteArrayList<Exception> exceptions = new CopyOnWriteArrayList<>();
    @Override
    public void onCompletion(final RecordMetadata metadata, final Exception exception) {
        if (exception != null) {
            exceptions.add(exception);
        }
    }
    public void check() throws Exception {
        if (!exceptions.isEmpty()) {
            // 聚合所有异常抛出,也可以根据自身需求调整处理逻辑
            Exception firstEx = exceptions.get(0);
            for (int i = 1; i < exceptions.size(); i++) {
                firstEx.addSuppressed(exceptions.get(i));
            }
            throw firstEx;
        }
    }
}

内容的提问来源于stack exchange,提问作者Karsten Spang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 17:24:00