Kafka消息发送异常排查:闲置后重试报错重启恢复
问题:Java Kafka客户端闲置数小时后发送消息失败,重启后恢复
问题描述
我们有一个Java进程向Confluent托管的Kafka发送消息,当进程闲置数小时后重试发送时,所有消息均触发报错,但重启进程后即可恢复正常。原本计划查看Broker日志但暂无权限,请问从客户端侧可采取哪些措施进一步调试?
客户端报错日志
2023-06-02 17:02:20,632 pid:[id] ERROR kafka.SimpleProducer - java.util.concurrent.ExecutionException: org.apache.kafka.common.InvalidRecordException: This record has failed the validation on broker and hence be rejected. at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.valueOrError(FutureRecordMetadata.java:98) at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.get(FutureRecordMetadata.java:67) at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.get(FutureRecordMetadata.java:30)
后续获取到的Broker日志
{"exception":{"stacktrace":"org.apache.kafka.common.InvalidRecordException: One or more records have been rejected due to 1 record errors in total, and only showing the first three errors at most: [RecordError(batchIndex=0, message='Log record DefaultRecord(offset=0, timestamp=1685667829424, key=19 bytes, value=2083 bytes) is rejected by the record interceptor io.confluent.kafka.schemaregistry.validator.RecordSchemaValidator')]","exception_class":"org.apache.kafka.common.InvalidRecordException","exception_message":"One or more records have been rejected due to 1 record errors in total, and only showing the first three errors at most: [RecordError(batchIndex=0, message='Log record DefaultRecord(offset=0, timestamp=1685667829424, key=19 bytes, value=2083 bytes) is rejected by the record interceptor io.confluent.kafka.schemaregistry.validator.RecordSchemaValidator')]"},"source_host":"xyz20092123456.obm.xxx.yyy.domain.com","method":"error","level":"ERROR","message":"[ReplicaManager broker=1] Error processing append operation on partition my-service-request-dev-xxx-2","mdc":{"brokerId":"1"},"@timestamp":"2023-06-02T01:03:49.490Z","file":"Logging.scala","line_number":"76","thread_name":"data-plane-kafka-request-handler-7","@version":1,"logger_name":"kafka.server.ReplicaManager","class":"kafka.utils.Logging"}
客户端侧调试措施
- 检查Schema Registry客户端缓存配置:客户端闲置期间缓存的Schema可能过期或失效,导致发送时使用的Schema与Broker验证器要求不匹配。查看
schema.registry.cache.size和schema.registry.cache.ttl参数,确认缓存是否会在闲置数小时后失效;调试阶段可临时禁用缓存,验证是否解决问题。 - 提升Schema相关日志级别:将
io.confluent.kafka.schemaregistry.client包的日志级别调整为DEBUG,记录闲置后首次重试时Schema的获取、验证过程,排查是否存在Schema获取失败、版本不匹配的情况。 - 记录发送时的Schema信息:在客户端代码中添加日志,输出每次发送消息使用的Schema ID、版本号,对比进程启动时和闲置后的Schema信息,确认是否有Schema在闲置期间被更新但客户端未同步。
- 检查Schema Registry连接状态:闲置数小时后,客户端与Schema Registry的HTTP连接可能已断开,导致无法正常验证Schema。查看连接池配置(如
http.max.idle.ms),调整参数保持长连接,或在发送前主动检查并重建连接。 - 模拟闲置场景抓包分析:在测试环境模拟进程闲置数小时的场景,抓取客户端与Kafka Broker、Schema Registry之间的网络请求,分析闲置后首次发送的请求内容,确认Schema是否正确携带、请求是否存在异常。
- 核查Producer连接配置:重点检查
connections.max.idle.ms(Kafka连接闲置超时)、retries、retry.backoff.ms等参数,确认闲置后Producer是否重新建立了与Broker的连接,重试请求是否存在异常。 - 添加本地Schema预验证:在发送消息前,调用Schema Registry客户端的验证接口对消息进行预验证,提前捕获Schema不匹配问题并记录详细错误,避免直接发送到Broker报错。
内容的提问来源于stack exchange,提问作者asb
相关产品推荐
相关产品推荐

