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

Kafka Connect S3源连接器忽略keys文件 如何配置同步消息key

问题根因

Confluent S3 Source连接器默认只读取存储消息value的数据文件,未显式配置key解析规则时,会自动过滤独立存储的.keys.json后缀key文件,这就是你日志中看到该类文件被移出扫描列表的直接原因。
另外你贴出的配置存在两处显性错误,会直接影响功能:

  1. 最后一行"transforms.AddPrefix.replacement":"$0_copy 缺失闭合双引号,属于JSON语法错误,连接器提交会直接失败
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 12:03:17