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

Kafka Connect Sink运行正常但sinkRecords始终为空求助

排查Kafka Connect Sink任务put()方法始终接收空集合的问题

问题场景

运行Kafka Connect服务,通过POST /connectors添加Sink连接器,配置如下:

{
  "name": "my-sink",
  "config": {
    "connector.class": "com.streamprocessing.MySinkConnector",
    "topics.regex": "MyComponent\\.MyTopic.*"
  }
}

连接器状态为RUNNING,任务id 0正常启动,但MySinkTask的put()方法每次接收的sinkRecords集合为空,日志每10秒打印Start put, batch size: 0。已确认:

  • 与Kafka Broker连接正常,偏移量/状态主题已创建
  • 连接器配置无问题,put()方法能被正常调用
  • 尝试过topics.regex=".*"或指定具体topics,均无效
  • 单独用消费者能正常消费测试消息,本地/测试环境现象一致

排查方向

1. 检查连接器消费者组的偏移量位置

Kafka Connect为每个连接器创建独立消费者组,命名规则为connect-${connector-name}(此处为connect-my-sink)。用Kafka命令行工具查看该组偏移量:

kafka-consumer-groups.sh --bootstrap-server <broker地址> --describe --group connect-my-sink

对比各主题分区的CURRENT-OFFSET与LOG-END-OFFSET:

  • 若两者相等,说明连接器已消费到主题末尾,无新消息可拉取。可重置偏移量到最早位置后重启连接器:
    kafka-consumer-groups.sh --bootstrap-server <broker地址> --reset-offsets --group connect-my-sink --to-earliest --execute --topic <目标主题>
    

2. 确认SinkConnector的任务配置传递是否正确

自定义MySinkConnector需在taskConfigs()方法中正确传递主题配置给SinkTask,检查该方法是否返回包含主题参数的配置:

@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
    List<Map<String, String>> configs = new ArrayList<>();
    Map<String, String> config = new HashMap<>(this.config);
    // 确保主题配置传递给Task
    config.put("topics.regex", config.get("topics.regex"));
    for (int i = 0; i < maxTasks; i++) {
        configs.add(config);
    }
    return configs;
}

若Task未收到主题配置,会导致无法订阅目标主题,自然拉不到消息。

3. 检查消息转换器配置是否匹配消息格式

Kafka Connect默认使用JsonConverter,若测试消息为非JSON格式(如字符串、二进制),需显式配置对应转换器:

{
  "name": "my-sink",
  "config": {
    "connector.class": "com.streamprocessing.MySinkConnector",
    "topics.regex": "MyComponent\\.MyTopic.*",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter"
  }
}

若转换器无法解析消息,可能导致消息被过滤,最终传递到put()的集合为空。可查看Kafka Connect日志,确认是否有反序列化错误。

4. 验证Kafka Connect的消费者拉取配置

检查全局或连接器级别的消费者配置:

  • consumer.max.poll.records:默认500,若设置过小且消息量极少,可能偶尔拉空,但不会持续每10秒空
  • poll.timeout.ms:默认3000ms,连接器会每隔该时间调用put(),即使无消息。你看到的10秒日志可能来自任务线程调度间隔或自定义逻辑
  • 可尝试在连接器配置中添加:
    "consumer.max.poll.records": "100",
    "consumer.poll.timeout.ms": "10000"
    
    调整后观察是否能拉到消息。

5. 检查SinkTask的start()方法是否干扰订阅逻辑

自定义MySinkTask的start()方法需避免覆盖默认订阅逻辑,仅初始化自定义资源:

@Override
public void start(Map<String, String> props) {
    // 仅初始化MongoDB连接等自定义资源,不要修改订阅逻辑
    this.mongoClient = new MongoClient(props.get("mongodb.uri"));
}

若start()方法存在错误的订阅操作,会导致Task无法正确订阅目标主题。

6. 确认主题分区与任务数的匹配

若目标主题有多个分区,可尝试将tasks.max设置为与分区数一致,验证是否能拉到消息:

{
  "name": "my-sink",
  "config": {
    "connector.class": "com.streamprocessing.MySinkConnector",
    "topics.regex": "MyComponent\\.MyTopic.*",
    "tasks.max": "3" // 假设主题有3个分区
  }
}

内容的提问来源于stack exchange,提问作者Yasin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 12:23:21