Kafka Connect中Camel S3 Sink的keyName属性异常问题求助
解决Camel S3 Sink连接器keyName属性不生效的问题
核心原因
Camel的AWS S3 Sink Kamelet默认启用分区策略,该机制会自动生成对象键,直接覆盖你配置的keyName属性。
具体解决步骤
1. 禁用自动分区
在连接器配置中添加参数,关闭默认分区行为:
"camel.kamelet.aws-s3-sink.partitioned": false
2. 用表达式语言引用Kafka消息键
要基于Kafka消息键生成S3对象键,使用Camel Simple表达式语法直接引用消息键:
"camel.kamelet.aws-s3-sink.keyName": "${header[kafka.KEY]}"
如果消息键是二进制格式(比如Debezium场景),需转换为字符串:
"camel.kamelet.aws-s3-sink.keyName": "${header[kafka.KEY].toString()}"
3. 完整示例配置
修改后的配置如下:
{ "connector.class": "org.apache.camel.kafkaconnector.awss3sink.CamelAwss3sinkSinkConnector", "camel.kamelet.aws-s3-sink.overrideEndpoint": true, "camel.kamelet.aws-s3-sink.uriEndpointOverride": "http://localstack:4566", "camel.kamelet.aws-s3-sink.bucketNameOrArn": "outbox", "camel.kamelet.aws-s3-sink.region": "eu-west-1", "camel.kamelet.aws-s3-sink.accessKey": "test", "camel.kamelet.aws-s3-sink.secretKey": "test", "camel.kamelet.aws-s3-sink.keyName": "${header[kafka.KEY]}", "camel.kamelet.aws-s3-sink.partitioned": false, "topics": "outbox.s3", "value.converter": "io.debezium.converters.BinaryDataConverter" }
额外注意事项
- 若需自定义复杂键格式(如添加时间戳前缀),可组合表达式:
"camel.kamelet.aws-s3-sink.keyName": "${date:now:yyyyMMdd}/${header[kafka.KEY]}" - 确保Kafka消息键非空,否则会 fallback 到自动生成的键值。
内容的提问来源于stack exchange,提问作者Boris Faniuk
相关产品推荐
相关产品推荐

