Kafka Connect因Confluent拦截器跳过消息的解决方案问询
解决Confluent监控拦截器跳过无时间戳消息的问题
问题概述
Kafka Connect日志中持续出现以下警告:
WARN Monitoring Interceptor skipped 2294 messages with missing or invalid timestamps for topic TEST_TOPIC_1. The messages were either corrupted or using an older message format. Please verify that all your producers support timestamped messages and that your brokers and topics are all configured with log.message.format.version, and message.format.version >= 0.10.0 respectively. You may also experience this if you are consuming older messages produced to Kafka prior to any of those changes taking place. (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor)
当前环境配置:
- Broker已设置:
KAFKA_INTER_BROKER_PROTOCOL_VERSION: 0.11.0 KAFKA_LOG_MESSAGE_FORMAT_VERSION: 0.11.0 - Kafka Connect保留了Confluent监控拦截器:
CONNECT_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" CONNECT_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" - 使用Pepperbox生产消息,消息未携带Kafka要求的时间戳元数据。
警告原因
Confluent监控拦截器依赖Kafka 0.10.0及以上版本的消息格式,该格式要求消息必须携带时间戳元数据(并非消息体中的字段)。Pepperbox默认生产的消息未设置该元数据,导致拦截器判定消息格式不符合要求,从而跳过统计。
解决方案
1. 配置Pepperbox自动添加消息时间戳
Pepperbox基于Kafka生产者客户端实现,只需为其添加生产者配置参数,让客户端自动为消息附加时间戳:
- 在Pepperbox的Kafka生产者配置中,添加
timestamp.type=CREATE_TIME参数:- 若使用JMeter图形界面:找到Pepperbox Sampler的「Producer Properties」配置项,新增键值对
timestamp.type->CREATE_TIME。 - 若使用配置文件:直接在生产者属性列表中加入该参数。
- 若使用JMeter图形界面:找到Pepperbox Sampler的「Producer Properties」配置项,新增键值对
该配置会让Kafka生产者自动将消息的创建时间作为时间戳元数据写入消息,满足拦截器的格式要求。
2. 验证配置生效
使用Kafka命令行工具验证消息是否已携带时间戳:
kafka-console-consumer.sh --bootstrap-server <你的Broker地址> --topic TEST_TOPIC_1 --from-beginning --property print.timestamp=true
若输出中包含CreateTime:前缀的时间戳,则说明配置生效,拦截器将不再跳过这些消息。
关于吞吐量的说明
自动添加时间戳属于轻量级操作,对生产者吞吐量的影响可以忽略不计,无需过度担心,可通过负载测试验证实际表现。
内容的提问来源于stack exchange,提问作者JayPatel
相关产品推荐
相关产品推荐

