Kafka Connect S3 Sink连接器无法消费指定Topic问题排查
我有一个从Kafka Topic读取数据并写入S3的S3 Sink连接器,但该连接器无法消费Topic中的数据。
连接器配置
{ "name": "my-connector", "config": { "connector.class": "io.confluent.connect.s3.S3SinkConnector", "behavior.on.null.values": "ignore", "s3.region": "<aws region>", "topics.dir": "my-topic", "flush.size": "1000", "tasks.max": "1", "s3.part.size": "5242880", "timezone": "UTC", "rotate.interval.ms": "30000", "locale": "en-US", "format.class": "io.confluent.connect.s3.format.avro.AvroFormat", "aws.access.key.id": "<aws access key>", "errors.deadletterqueue.topic.replication.factor": "1", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "s3.bucket.name": "<aws bucket>", "partition.duration.ms": "30000", "schema.compatibility": "NONE", "topics": "my-topic", "aws.secret.access.key": "<aws secret key>", "task.class": "io.confluent.connect.s3.S3SinkTask", "errors.deadletterqueue.topic.name": "dlq-my-topic", "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner", "name": "my-connector", "errors.tolerance": "all", "storage.class": "io.confluent.connect.s3.storage.S3Storage", "path.format": "'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH", "rotate.schedule.interval.ms": "60000", "timestamp.extractor": "Record" } }
连接器状态
执行命令:
curl localhost:8083/connectors/my-connector/status
返回结果:
{ "name": "my-connector", "connector": { "state": "RUNNING", "worker_id": "localhost:8083" }, "tasks": [ { "id": 0, "state": "RUNNING", "worker_id": "localhost:8083" } ], "type": "sink" }
Kafka消费者组信息
执行命令:
./kafka-consumer-groups.sh --bootstrap-server kafka:9092 --group connect-my-connector --describe
返回结果:
| GROUP | TOPIC | PARTITION | CURRENT-OFFSET | LOG-END-OFFSET | LAG | CONSUMER-ID | HOST | CLIENT-ID |
|---|---|---|---|---|---|---|---|---|
| connect-my-connector | my-topic | 0 | 1182 | 12072 | 10890 | connector-consumer-my-connector-0-68793e0d-8312-4d20-b23c-5221ca54b0dc | ip | connector-consumer-my-connector-0 |
可见Kafka Connect在消费者组中有活跃的消费者,但连接器仍无法消费Kafka数据。请问可能的原因是什么?参考过相关问题,但由于我们的Kafka集群未启用认证,该答案不适用。除认证外,还有哪些原因会导致已连接的Kafka Connect连接器无法消费Topic消息?启动、插件加载及连接器初始化阶段日志均无明显错误。
1. 格式转换器与输出格式不兼容
你的配置中,value.converter使用org.apache.kafka.connect.json.JsonConverter(将Kafka消息转为JSON结构),但format.class配置的是io.confluent.connect.s3.format.avro.AvroFormat(要求写入Avro格式数据)。这种不匹配会导致消息处理失败:
- 若
JsonConverter输出无Schema的JSON,AvroFormat无法解析;即使是带Schema的JSON,也不符合Avro格式要求。 - 由于配置了
errors.tolerance: "all",错误会被容忍,但消息无法正常处理,偏移量也无法推进。
排查建议:
- 将
value.converter改为io.confluent.connect.avro.AvroConverter(配合AvroFormat使用),或者将format.class改为io.confluent.connect.s3.format.json.JsonFormat,保持转换器与输出格式一致。
2. 时间戳提取器配置异常
你设置了timestamp.extractor: "Record",意味着连接器会从消息内容中提取时间戳进行时间分区。如果出现以下情况,会导致消息处理停滞:
- 消息中不存在连接器预期的时间戳字段(默认字段为
timestamp,可通过timestamp.field配置修改); - 时间戳字段格式不符合要求(比如不是长整型的毫秒时间戳)。
排查建议:
- 检查消息结构,确认包含正确的时间戳字段;
- 临时将
timestamp.extractor改为"Wallclock"(使用系统当前时间),验证是否能正常消费消息。
3. JsonConverter的Schema配置缺失
JsonConverter默认要求消息携带Schema信息(schemas.enable默认值为true)。如果你的Kafka消息是纯JSON(无Schema),转换器会抛出转换错误,进而导致消息处理卡住。
排查建议:
- 在连接器配置中添加
"value.converter.schemas.enable": "false",允许处理无Schema的JSON消息。
4. 分区器逻辑阻塞
使用TimeBasedPartitioner时,如果时间戳提取失败,分区器无法确定消息的写入路径,会导致整个处理流程阻塞,无法推进偏移量。这通常和第2点的时间戳提取问题关联。
排查建议:
- 先解决时间戳提取问题,再验证分区器是否正常工作;
- 可以临时改用
FieldPartitioner或DefaultPartitioner,排除分区器的影响。
5. Kafka Connect Worker资源瓶颈
如果Worker节点的CPU、内存资源不足,会导致消息处理线程被阻塞,无法正常消费和处理消息:
- JVM内存不足引发频繁GC,导致线程停顿;
- CPU使用率过高,处理速度跟不上消息生产速度。
排查建议:
- 检查Worker的JVM日志,查看是否有GC频繁、内存溢出的记录;
- 查看系统资源监控,确认CPU、内存使用率是否在合理范围。
6. DLQ Topic异常
虽然配置了errors.tolerance: "all"和DLQ,但如果DLQ Topic不存在或者无法写入(比如Topic未创建、副本因子不匹配),可能会导致连接器处理消息时卡住。
排查建议:
- 检查
dlq-my-topic是否已创建,且副本因子与配置的errors.deadletterqueue.topic.replication.factor一致; - 查看Worker日志中是否有DLQ相关的错误信息。
内容的提问来源于stack exchange,提问作者Mr T.

