如何在Kafka SFTP CSV连接器中设置Null Key实现轮询分区分配
问题描述
我在用Kafka的SFTP CSV Source Connector时,当前配置生成的Key是空结构体{},导致所有消息都挤在单个分区里。根据文档,设置Null Key能让消息在Broker间轮询分配——我用Python脚本不指定Key的时候确实能达到这个效果,现在想知道怎么通过连接器配置实现同样的结果。
我的连接器配置如下:
{ "name": "NdpiSourceConnector_csv_1", "config": { "tasks.max": "1", "connector.class": "io.confluent.connect.sftp.SftpCsvSourceConnector", "cleanup.policy":"MOVE", "behavior.on.error":"IGNORE", "input.path": "/home/******/Desktop/ndpi/drop", "error.path": "/home/******/Desktop/ndpi/error", "finished.path": "/home/shyena/Desktop/ndpi/processed", "input.file.pattern": ".*\\.csv", "sftp.username":"********", "sftp.password":"**********", "sftp.host":"**************", "sftp.port":22, "kafka.topic": "NdpiSourceTopic1", "csv.first.row.as.header": "true", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable":"false", "key.schema":"{\"type\":\"STRUCT\",\"name\":\"FlowRecord\",\"isOptional\":\"true\",\"fieldSchemas\":{}}", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.schema": "{\"type\":\"STRUCT\",\"name\":\"FlowRecord\",\"fieldSchemas\":{\"flow_id\":{\"type\":\"STRING\"},\"protocol\":{\"type\":\"STRING\"},\"first_seen\":{\"type\":\"STRING\"},\"last_seen\":{\"type\":\"STRING\"},\"duration\":{\"type\":\"STRING\"},\"src_ip\":{\"type\":\"STRING\"},\"src_port\":{\"type\":\"STRING\"},\"dst_ip\":{\"type\":\"STRING\"},\"dst_port\":{\"type\":\"STRING\"}}}", "value.converter.schemas.enable":"false" } }
解决方案
要让连接器生成Null Key,实现消息轮询分配到不同分区,有两种简单的配置修改方式:
方式一:用NullConverter作为Key转换器
直接把key.converter改成org.apache.kafka.connect.converters.NullConverter,同时删掉所有多余的Key相关配置(比如key.converter.schemas.enable和key.schema)。修改后的关键配置片段:
"key.converter": "org.apache.kafka.connect.converters.NullConverter"
这个转换器就是专门用来生成Null Key的,是最直接的解决办法。
方式二:移除Key配置让连接器自动生成Null Key
如果不想换转换器,可以删除key.schema这一行,同时去掉key.converter.schemas.enable配置。SFTP CSV Source Connector在没配置Key生成规则的情况下,默认会发送Null Key的消息。修改后这部分配置只保留:
"key.converter": "org.apache.kafka.connect.json.JsonConverter"
验证方法
改完配置重启连接器后,你可以用以下命令验证效果:
- 查看主题分区情况:
kafka-topics.sh --describe --topic NdpiSourceTopic1 --bootstrap-server <你的Broker地址>
- 消费消息并查看分区分配:
kafka-console-consumer.sh --topic NdpiSourceTopic1 --bootstrap-server <你的Broker地址> --property print.partition=true
如果消息分散在不同分区,说明配置生效了。
内容的提问来源于stack exchange,提问作者NOOBCODER
相关产品推荐
相关产品推荐

