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
相关产品推荐
相关产品推荐

