Kafka Connect S3源连接器忽略keys文件 如何配置同步消息key
问题根因
Confluent S3 Source连接器默认只读取存储消息value的数据文件,未显式配置key解析规则时,会自动过滤独立存储的.keys.json后缀key文件,这就是你日志中看到该类文件被移出扫描列表的直接原因。
另外你贴出的配置存在两处显性错误,会直接影响功能:
- 最后一行
"transforms.AddPrefix.replacement":"$0_copy缺失闭合双引号,属于JSON语法错误,连接器提交会直接失败 topics.dir参数配置错误,该参数指桶内的相对存储路径,不需要重复填写桶名,你当前写的test-bucket/topics会导致连接器拼接出test-bucket/test-bucket/topics/...的错误S3路径,根本扫不到目标文件
修复方案
在配置中补充key文件的格式声明,同时修正上述两处错误即可,不需要改动其他分区规则、路径规则,连接器会自动匹配同偏移量下的key文件和value文件,组装成完整消息写入Kafka:
- 新增
key.format.class参数,值和value使用的格式类保持一致,指定key文件的解析格式
修正后的完整配置
{ "connector.class": "io.confluent.connect.s3.source.S3SourceConnector", "tasks.max": "1", "s3.region": "eu-central-1", "s3.bucket.name": "test-bucket", "topics.dir": "topics", "format.class": "io.confluent.connect.s3.format.json.JsonFormat", "key.format.class": "io.confluent.connect.s3.format.json.JsonFormat", "partitioner.class": "io.confluent.connect.storage.partitioner.DefaultPartitioner", "schema.compatibility": "NONE", "confluent.topic.bootstrap.servers": "blockchain-kafka-kafka-0.blockchain-kafka-kafka-headless.default.svc.cluster.local:9092", "transforms": "AddPrefix", "transforms.AddPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.AddPrefix.regex": ".*", "transforms.AddPrefix.replacement": "$0_copy" }
校验标准
配置更新生效后,连接器日志中不会再出现.keys.json文件被移除的调试记录,消费对应*_copy后缀的目标主题时,可同时读取到消息的key和value内容。
内容的提问来源于stack exchange,提问作者Yordan Pavlov
相关产品推荐
相关产品推荐

