Kafka Connect Sink运行正常但sinkRecords始终为空求助
问题场景
运行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

