调用Kafka Producer返回的Future<RecordMetaData>如何判断执行成功或失败?
Kafka Producer 异步发送结果判断规则
首先明确结论:仅拿到Future<RecordMetadata>对象不等于消息发送成功。
Kafka Producer默认是异步发送模式,调用send()方法后,消息会先存入客户端的本地缓冲区,只要消息能正常入队(没有出现序列化失败、缓冲区满且阻塞超时、配置非法等前置校验错误),就会直接返回Future对象,不会抛出异常,此时消息还没有发往Broker,也没有得到Broker的任何确认。
你可以通过两种方式判断最终发送是否成功:
- 同步阻塞获取结果
调用Future.get()方法阻塞当前线程,直到发送流程结束:- 方法返回
RecordMetadata实例:发送成功,实例中包含消息对应的主题、分区、偏移量、写入时间戳等元信息 - 方法抛出异常:发送失败,可捕获异常类型判断具体失败原因,常见的有
TimeoutException(请求超时)、SerializationException(消息序列化失败)、AuthorizationException(权限不足)等
- 方法返回
- 异步回调(推荐非阻塞场景使用)
调用send()方法时传入Callback实现,发送结果会通过回调异步通知,不需要阻塞等待:producer.send(producerRecord, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception == null) { // 发送成功,处理元数据 } else { // 发送失败,处理异常 } } });
注意:如果你配置了
retries参数大于0,可重试类型的异常会先触发客户端自动重试,重试次数耗尽后才会判定为失败;如果配置acks=0,消息只要写入本地缓冲区就会标记为成功,不会等待Broker确认,这种情况下即使消息在网络传输中丢失,你也不会收到报错。
内容的提问来源于stack exchange,提问作者stackerstack
相关产品推荐
相关产品推荐

