尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spring Kafka发送回调机制:ListenableFutureCallback与ProducerListener实战解析

Spring Kafka发送回调机制:ListenableFutureCallback与ProducerListener实战解析 写 Kafka 的兄弟十有八九都遇到过这种情况kafkaTemplate.send()这行代码一写日志也不报错消费者那边就是收不到消息或者消息内容对不上。等排查到最后才猛然发现问题出在发送结果上——你压根没确认消息到底有没有发出去或者发是发出去了但 broker 端写入失败被静默吞掉了。所以 Spring Kafka 里这套发送回调机制不是我吓唬人它就是你消息系统的最后一道安全网。这篇东西我打算把ListenableFutureCallback和ProducerListener这两个机制一次性讲透。前者让你为单条消息注册异步回调后者让你在生产端做全局监听和指标统计两个配合起来基本能做到每条消息的去向可查、每个失败点可追踪、每类异常可报警。适合正在接手 Kafka 生产环境的同学、被消息凭空消失坑过的人以及想在 Spring Boot 项目里把 Kafka 发送链路做成可观测的团队。1. 项目概述发送回调为什么绕不开1.1 一个真实的生产事故send() 成功不等于消息被接受先说一个我踩过的坑非常典型。有一次订单系统上线业务方反馈说部分订单没有进入后续处理流程单查数据库记录订单状态明明是已创建。我们查了半天日志发现kafkaTemplate.send()全部返回成功方法也没有抛任何异常但下游数据分析团队用的消息就是缺了一部分。后来定位到的原因很简单当时测试环境为了追求性能把acks配成了 0也就是 producer 只管把消息丢给 socket 缓冲区压根不关心 broker 是否真的写入成功。send()方法是异步的它返回一个Future这个 Future 只在表示发送动作已提交并不代表数据已经落盘或者被分区副本接收。你把acks0配上去网络抖动一下数据丢了都不知道。这次事故之后我给自己定了一条规矩凡是调用kafkaTemplate.send()的地方必须配套处理发送结果。要么用ListenableFutureCallback做业务兜底要么在全局配置ProducerListener做统一监控两条至少占一条。这就是这两个机制存在的根本意义——Kafka 本身是个分布式的、消息可以乱序可以丢失的系统你要是不关心发送结果就等于闭着眼睛往河里扔漂流瓶。1.2 两个机制一个管单点业务一个管全局观测很多人会把ListenableFutureCallback和ProducerListener搞混以为随便用哪个都行其实两者定位完全不同。ListenableFutureCallback是附着在单次send()调用上的它对应的是KafkaTemplate.send()返回的那个ListenableFuture。你在这个回调里能拿到SendResult里面塞着RecordMetadata可以精确到这条消息写入了哪个分区、对应 offset 是多少、服务端时间戳是什么。这玩意儿适合做跟业务强相关的处理比如发送失败后写补偿表、记录日志、触发重试或者把实际落盘的 offset 回传给下游做对账。ProducerListener则是挂在 producer 层面的全局监听器它相当于是给所有通过这个 producer 发送的消息装了一个监控探头。这个接口能监听到每一条消息的发送成功事件、失败事件不管你是在哪个业务代码里调的send()。通常拿它做统一埋点、指标统计、异常告警而不是做具体的业务补偿——因为它在全局层面它不知道也不关心你这业务是谁。用一句话总结ListenableFutureCallback是给单条消息请的贴身保镖ProducerListener是给整个生产端装的监控摄像头。两者不冲突可以同时存在实际项目中我推荐两个都用。2. 核心细节解析先吃透 KafkaTemplate 返回的 ListenableFuture2.1 KafkaTemplate.send() 背后到底发生了什么要搞明白回调机制得从send()方法的实现说起。Spring Kafka 的KafkaTemplate.send()底层调用的是KafkaProducer.send(record, callback)这个方法是标准 Kafka Client 里最核心的发送入口。KafkaProducer 内部维护了一批RecordBatch发送消息时先把消息追加到对应的批次里然后由一个叫做Sender的后台线程负责把批次数据真正写到 broker。注意关键词后台线程。也就是说send()方法本身是异步的调用之后立刻返回数据什么时候真正发出去由sender线程决定。你拿到的Future是 send 动作已经提交到内存缓冲区的凭据而不是broker 已确认的凭据。在acks0的情况下这个 Future 甚至不会等待任何服务端响应直接就算成功。KafkaProducer 的send()方法签名是这样的public FutureRecordMetadata send(ProducerRecordK, V record, Callback callback)那 KafkaProducer 自带的那个Callback是什么它和 Spring Kafka 的ListenableFutureCallback是什么关系简单说Spring Kafka 做了一层适配。它把 KafkaProducer 的异步发送封装成了一个可以被 Spring 容器管理的、符合异步语义的ListenableFuture同时把KafkaProducer.Callback内部转接成了 Spring 的ListenableFutureCallback。所以你给kafkaTemplate.send()注册回调的时候本质上就是在给 KafkaProducer 内部的 callback 注册回调。再去翻 Spring Kafka 的源码你会发现KafkaTemplate.doSend()内部有一段类似这样的逻辑ListenableFutureSendResultK, V future new SettableListenableFuture(); producer.send(new ProducerRecord(topic, partition, key, value), (metadata, exception) - { if (exception null) { future.set(new SendResult(record, metadata)); } else { future.setException(new KafkaProducerException(record, Send failed, exception)); } });这结构其实不难理解SettableListenableFuture是一个可手动填充结果和异常的 Future 实现等 KafkaProducer 的 callback 触发了再把结果或者异常塞进去。你调future.addCallback(...)注册的回调就会在结果或异常被设置的时候被触发。2.2 注册回调的两种姿势addCallback 与时下更推荐的 whenComplete知道send()返回的是ListenableFutureSendResultK, V之后编写回调就有几种选择。最传统的方式是调用ListenableFutureCallbackListenableFutureSendResultString, String future kafkaTemplate.send(order-event, order-123, created); future.addCallback(new ListenableFutureCallbackSendResultString, String() { Override public void onSuccess(SendResultString, String result) { RecordMetadata metadata result.getRecordMetadata(); log.info(消息发送成功, topic{}, partition{}, offset{}, metadata.topic(), metadata.partition(), metadata.offset()); } Override public void onFailure(Throwable ex) { log.error(消息发送失败, 原因: {}, ex.getMessage(), ex); // 这里可以做业务补偿 } });但说实话匿名内部类写在业务代码里有点丑尤其是你只是想做一段简单处理的时候。Spring Kafka 2.8 以后ListenableFuture扩展成了CompletableFuture的父接口准确说是引入了新的实现于是你可以用更现代的方式kafkaTemplate.send(order-event, order-123, created) .whenComplete((result, ex) - { if (ex null) { RecordMetadata metadata result.getRecordMetadata(); // success } else { // failure } });这种方法链读起来舒服很多也更容易配合 Lambda 表达式。而且CompletableFuture风格的 API 有丰富的后续方法比如thenAccept、exceptionally可以做链式编排。我个人的习惯是如果只是记录日志和简单重试用whenComplete如果要对 Future 做复杂编排、超时控制或者组合多个发送结果用CompletableFuture风格更顺手。还有一种偷懒的姿势是future.get()同步获取结果。这个方法会阻塞当前线程直到 broker 返回写入结果用起来简单但会严重拖慢吞吐。除非你是在单元测试里写断言否则我不建议在真实业务代码里用get()——一旦 broker 响应慢它会直接卡死你的业务线程。2.3 回调执行的线程模型这个坑必须提前避回调代码在哪个线程里执行这是一个非常关键、又很容易被忽视的问题。我需要明确告诉你ListenableFutureCallback的回调是在 KafkaProducer 的 sender 线程中执行的严格来说是 sender 线程回调 KafkaProducer.Callback 之后Spring 又在这个链路里触发了 Future 的结果设置。这意味着两件事第一回调代码不要做耗时操作。你如果在回调里写数据库操作、调远程接口、做复杂的逻辑处理就会占用 sender 线程。sender 线程是负责批量发送消息的它本身要轮询 broker 的响应、维护元数据、管理批次重试。你阻塞它就会连累整个生产者导致后面所有消息的发送延迟升高、吞吐下降。我见过有人把订单校验服务直接写在回调里上线后 Kafka 生产端延迟直接翻了好几倍就是这个原因。正确做法是回调里只做轻量处理需要重的逻辑就丢到线程池里去异步执行。第二回调里的异常处理一定要做。回调方法如果抛出未捕获异常在 KafkaProducer 内部它会被捕获并打印 WARN 日志具体取决于版本你的业务逻辑可能已经执行了一半了再往下就没人管了。所以不管哪种方式注册的回调onFailure/异常分支里务必要 try-catch 兜底至少把异常打出来避免排查问题时一头雾水。3. ProducerListener 详解全局消息发送监听器3.1 ProducerListener 接口结构与版本差异如果说ListenableFutureCallback是给单个发送点准备的那ProducerListener就是在整个 producer 生命周期上挂载的全局监听器。这个接口最早在 Spring Kafka 1.x 时代就有了但不同版本接口定义有差异这一点千万要注意。以 Spring Kafka 3.x 为例ProducerListener 的关键方法有这几个public interface ProducerListenerK, V { default void onSuccess(ProducerRecordK, V record, RecordMetadata metadata) { } default void onError(ProducerRecordK, V record, RecordMetadata metadata, Exception exception) { } default boolean isInterestedInSuccess() { return false; } }注意看默认方法全是空实现没有强制要求。onSuccess在消息发送成功后触发onError在发送失败时触发。还有一个容易被忽略的isInterestedInSuccess()——它的返回值决定了 Spring Kafka 要不要调用你的onSuccess方法。这个设计是有原因的成功消息的数量通常很大如果你只关心失败场景不关心成功场景就不希望系统为每一条成功消息都多一次额外的回调调用。所以isInterestedInSuccess()默认返回false表示我不关心成功回调返回true才开启成功回调。在 Spring Kafka 3.x 里这个方法的默认实现是false。很多初学的人会发现onSuccess怎么不触发问题就在这儿。另外在 Spring Kafka 2.x 的早期版本里ProducerListener 还有onAck这样的方法后来被重构掉了。所以你要是参考到很老的博客一定先确认一下版本别把不存在的接口方法套上去。3.2 手写一个带指标监控的 ProducerListener我自己在实际项目里维护了一个全局的 ProducerListener负责三点打日志、统计指标、失败告警。核心实现大概是这样的public class MonitoringProducerListenerK, V implements ProducerListenerK, V { private final MeterRegistry meterRegistry; private final KafkaProducerMonitor monitor; public MonitoringProducerListener(MeterRegistry meterRegistry, KafkaProducerMonitor monitor) { this.meterRegistry meterRegistry; this.monitor monitor; } Override public boolean isInterestedInSuccess() { // 开启成功回调因为我们需要统计成功率 return true; } Override public void onSuccess(ProducerRecordK, V record, RecordMetadata metadata) { monitor.recordSuccess(record.topic(), record.partition() ! null ? record.partition() : metadata.partition()); meterRegistry.counter(kafka.producer.message.success, topic, record.topic()) .increment(); } Override public void onError(ProducerRecordK, V record, RecordMetadata metadata, Exception exception) { String topic record.topic(); String exceptionClass exception.getClass().getSimpleName(); monitor.recordFailure(topic, exceptionClass); meterRegistry.counter(kafka.producer.message.failure, topic, topic, exception, exceptionClass) .increment(); log.error(Kafka producer 发送失败, topic{}, key{}, value{}, error{}, topic, record.key(), record.value(), exception.getMessage()); } }这段代码放在生产环境里是完全可用的。通过MeterRegistryMicrometer 的核心抽象把指标暴露给 Prometheus 或者 InfluxDB然后用 Grafana 画大盘。每次发送失败都能立刻在监控上体现出来报警规则设置成功率低于 99.99% 触发警告即可。有一个细节值得琢磨ProducerListener.onSuccess和ProducerListener.onError拿到的参数里ProducerRecord包含了完整的消息体包括 key、value。这意味着你可以在失败的时候记录下具体的消息内容。但注意隐私问题和日志量问题——失败消息的频率如果高把 value 全部打出来磁盘很快就撑爆了。我一般只记录 key 和 topicvalue 太大就不打了要查原始内容直接去源业务库里捞。3.3 在 Spring Boot 中装配与替换默认 ProducerListener那 ProducerListener 怎么挂到 Spring Boot 容器里这里有个特别容易踩的坑在 Spring Boot 自动配置下你直接往容器里丢一个 ProducerListener BeanKafkaTemplate 默认是不用的。原因在于 Spring Boot 的 KafkaAutoConfiguration 在创建KafkaTemplate的时候默认不会去容器里自动查找ProducerListener并注入。你要么在创建KafkaTemplate时手动设置要么直接自定义一个 ProducerFactory 的 Bean。我推荐的做法是自定义配置类显式创建一个带 ProducerListener 的KafkaTemplateConfiguration public class KafkaProducerConfig { Bean public ProducerFactoryString, String producerFactory(ProducerProperties properties) { MapString, Object props new HashMap(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, properties.getBootstrapServers()); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 其他参数 props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.RETRIES_CONFIG, 3); DefaultKafkaProducerFactoryString, String factory new DefaultKafkaProducerFactory(props); // 注意这里如果消息发送失败ProducerFactory 需要能拿到异常信息 factory.setProducerListener(new MonitoringProducerListener(meterRegistry, monitor)); return factory; } Bean public KafkaTemplateString, String kafkaTemplate(ProducerFactoryString, String producerFactory) { return new KafkaTemplate(producerFactory); } }比较关键的是DefaultKafkaProducerFactory.setProducerListener()这个方法。设置之后从该工厂创建出来的所有 producer 都会带上这个监听器后续你注入的KafkaTemplate发送的任何消息都会经过这个监听器。如果你用的是 Spring Boot 的自动配置不想自己定义整个 ProducerFactory那也能绕手动构建一个KafkaTemplateBean使用自动配置好的ProducerFactory然后调用kafkaTemplate.setProducerListener(...)。不过要注意顺序问题最好用ObjectProvider拿到自动配置的 factory避免 Bean 创建顺序导致的坑。4. 项目实操构建一套完整的发送确认与追踪方案4.1 第一步确认 ack 配置让回调有意义说句实在话如果你生产端的acks配置错了后面所有的回调都像是在瞎忙活。从语义上讲回调能拿到的成功结论其实完全取决于 ack 级别。所以动笔画代码之前先把这个配置搞清楚。acks0producer 不等待 broker 任何确认send()直接返回成功。消息发出去之后broker 是否真的写入、有没有报错调用方完全不知道。这种配置下回调的成功分支收到的是发送动作已完成毫无可靠性。acks1leader 分区写入成功就算成功不等待 follower 同步。大多数中低要求场景够用但 leader 挂掉尚未同步到 follower 时数据可能丢。acksall等价于acks-1等待所有 in-sync 副本ISR都写入成功才算成功可靠性最强。生产环境必须用这个再配合min.insync.replicas设置最小同步副本数。我们生产环境的标准配置是这样的spring: kafka: producer: acks: all retries: 3 properties: enable.idempotence: true max.in.flight.requests.per.connection: 5 delivery.timeout.ms: 120000 request.timeout.ms: 30000这里多说一句enable.idempotence开启幂等发送之后producer 会给每条消息加上一个递增的序列号broker 据此去重配合retries消息重复发送的概率会大幅度降低。幂等发送要求retries 0、acksall、max.in.flight.requests.per.connection 5新版默认放宽了这个限制。开了幂等之后再配合回调确认消息可靠性基本稳了。4.2 第二步用 SendResult 做业务补偿和重试配置到位之后回调里的信息就有了真实可靠的含义。SendResult里最有价值的东西是RecordMetadata它能告诉你消息最终落在了哪个 physical partition、写到什么 offset、broker 时间戳是多少。这些数据对于订单这类业务来说就能用于和下游对账。我做过一个订单事件补偿方案思路是这样的Transactional public void sendOrderEvent(OrderEvent event) { // 1. 先把订单事件持久化到本地事件表状态为 PENDING orderEventRepository.save(OrderEventRecord.pending(event)); // 2. 发送消息在回调里根据结果更新状态 ListenableFutureSendResultString, String future kafkaTemplate.send( order-event, event.getOrderId(), objectMapper.writeValueAsString(event)); future.whenComplete((result, ex) - { if (ex null) { // 更新本地事件表状态为 SENT并记录 partition/offset orderEventRepository.markSent(event.getEventId(), result.getRecordMetadata().partition(), result.getRecordMetadata().offset()); } else { // 记录失败定时任务扫描并重试 log.warn(订单事件发送失败, eventId{}, 进入重试队列, event.getEventId(), ex); orderEventRepository.markFailed(event.getEventId(), ex.getMessage()); } }); }配合一个定时任务每分钟扫描一次本地事件表把状态为FAILED或超时仍未更新的记录捞出来重新发送。重试超过 3 次就进入人工处理队列。这套方案实现简单却能应对大多数发送失败场景。要注意的一个坑是回调里的数据库写入不能放在和send()同一个事务里。KafkaTemplate.send()是异步的回调可能发生在事务提交之后也可能发生在事务提交之前你跟事务放一起时序会很混乱。正确姿势是把事件记录保存在事务里发送动作在事务外触发回调更新状态不要依赖原事务的上下文。4.3 第三步把 ProducerListener 接入监控大盘回调做业务兜底是不够的你还得有全局视野。我前文提到了用 Micrometer 统计指标这里再说细一点给出一个可以抄作业的配置。引入 Spring Boot Actuator 和 Micrometer 之后只需要在 application.yml 里加上management: endpoints: web: exposure: include: health,info,prometheus,metrics metrics: export: prometheus: enabled: true接着把 ProducerListener 里的指标统计用 Micrometer 的 Tag 组织好对应 Prometheus 拉取之后我一般会在 Grafana 面板上显示三个核心监控项生产端发送成功率按 topic 分组5 分钟内成功率低于 99.99% 就报警。发送失败原因 Top 排行榜按异常类型聚合能快速看出是TimeoutException、KafkaException还是AuthorizationException。回调耗时和积压量在回调里记录处理耗时如果耗时突然升高很可能回调里写了重逻辑需要检查。这个方案部署之后整条 Kafka 发送链路的状态都是可见的。哪条 topic 出问题、哪个分区出问题、最近什么时候开始出现超时一目了然。再配合可视化管理工具查 broker 状态和消费组 offset基本能覆盖生产环境 90% 的消息类故障排查。5. 常见问题与排查技巧实录5.1 回调方法没执行先排查这五个点阶段性地总会有人来问我代码明明写了回调为什么不触发啊。我总结了五个最常见的排查方向第一acks0 会导致失败回调不触发。acks0时producer 压根不关心 broker 的结果KafkaProducer 内部直接走成功分支所以即便消息实际没写进去你的onFailure也永远不会跑。遇到没报错但数据丢了的情况第一个看 ack 配置。第二isInterestedInSuccess()返回值导致成功回调被跳过。这个问题在 ProducerListener 上尤其常见。默认返回falseonSuccess不会执行。只有把返回值改成trueSpring Kafka 才会把成功通知传给你。第三你设置 ProducerListener 的对象不对。如果你的KafkaTemplate没有使用你配置的 ProducerFactory或者你用了两个不同的 ProducerFactory 但只在其中一个上设置了监听器就会出现某个 Topic 有监控、某个 Topic 没监控的现象。排查时先确认KafkaTemplate的producerFactory到底是谁。第四回调里抛异常导致调用链被吞掉。这个问题前面提过回调方法里如果有未捕获异常KafkaProducer 只会在内部记录一下你的后续逻辑就断了。建议回调里 try-catch 打日志至少把现场保存下来。第五Application 提前退出导致回调没有执行完。如果你的程序是用main方法直接跑的发送消息后没等回调执行程序就结束了回调自然就不会触发。Spring Boot 场景一般不会遇到但批量任务、单元测试里很常见。5.2 报 cluster authorization failed 等典型发送异常速查Kafka 发送阶段的报错五花八门我挑几种高频的列个速查表方便大家对照异常信息可能原因排查方向cluster authorization failed客户端没有集群级权限常见于大集群多租户环境检查 ACL 授权确认是否允许CLUSTER_ACTIONtopic authorization failed当前用户没有该 topic 的写权限为业务账号配置对应的 topic 写入权限TimeoutException: Topic xxx not present in metadata after 60000 mstopic 不存在或 broker 元数据拉取异常用客户端/repl 查看 topic确认 bootstrap-server 地址正确RecordTooLargeException消息体大小超过 max.request.size增大 max.request.size 或调整消息拆分策略KafkaException: Could not obtain metadata元数据拉取失败网络不通或 broker 未就绪检查网络、broker 状态确认 DNS 可解析BufferExhaustedException发送缓冲区和积压量过大内存队列满了调大 buffer.memory 或加快消费侧吞吐处理这些报错有一个通用原则先把 Kafka producer 控制台输出的全部日志打开再看 broker 端的日志最后结合监控大盘确定是客户端问题还是服务端问题。很多时候发送端日志的告警级别默认比较高看不到细节建议把 Kafka client 的日志级别调到 DEBUG特别是org.apache.kafka.clients.producer.internals包下的排错会快很多。5.3 消息延迟高回调与批量发送的相互干扰最后一个常见问题是生产端消息延迟居高不下。很多人的第一反应是网络慢、broker 慢但其实回调和批量发送的配置也会显著影响延迟。先说一个真实的坑。有个系统给 KafkaTemplate 的每条消息都加了一个ListenableFutureCallback回调里做了一次try { Thread.sleep(50); } catch ...模拟耗时处理。表面上业务代码没受影响但前面我们说过回调是在 sender 线程里执行的。50 毫秒看起来不长但对于高吞吐的 sender 线程影响反映在大量消息的整体延迟上。我们调监控发现 producer 端 batch 迟迟攒不满因为 sender 线程被回调卡住了后续消息排队积压。解决方式是把回调里的耗时逻辑丢到独立的线程池sender 线程只负责标记结果和触发回调方法。另外延迟高还可能跟这几个参数有关linger.ms如果配得太大消息会在本地缓冲区多等一段时间延迟自然高batch.size配置过小时批次频繁刷出浪费网络带宽同样会引起吞吐下降和响应时间波动max.in.flight.requests.per.connection影响 broker 在途请求数如果配得太大一旦 broker 处理不过来排队反而变长。调试时别一次性改一堆参数我建议只动一个变量收集指标再改下一个。重点关注kafka.producer.request.latency.avg和kafka.producer.batch.size.avg这类指标能直观反映调整效果。最后再补充一个我个人实际使用的心得这套回调机制我用下来最大的体会是它既不能替代消息一致性方案也不能替代监控体系但它是两者之间的重要桥梁。ListenableFutureCallback的粒度太细不适合做全局观测ProducerListener的粒度太粗不适合做业务补偿。把两者组合起来一个在面对单条消息时能精确补偿一个在系统层面能快速发现趋势异常才算把这个发送确认的逻辑真正闭环了。还有一个小技巧分享给大家如果线上环境允许给你的ProducerListener增加一个开关通过配置中心动态开启或关闭成功回调。平时只保留失败监控排查问题的时候临时打开成功回调多收几分钟数据就能确认 消息是否全部发出了、各分区分布是否均衡。开关功能很简单就是一个AtomicBoolean的判断但排查问题时非常管用。这块现在已经成为我们团队处理消息类问题时的标准操作了。
返回列表