You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

返回结果:

GROUPTOPICPARTITIONCURRENT-OFFSETLOG-END-OFFSETLAGCONSUMER-IDHOSTCLIENT-ID
connect-my-connectormy-topic011821207210890connector-consumer-my-connector-0-68793e0d-8312-4d20-b23c-5221ca54b0dcipconnector-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.

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 22:04:54