使用Apache Camel Source实现S3到Kafka同步的配置问题咨询
需求1:不删除源文件且保证消费不重复
可以仅通过配置实现,你目前配置的camel.source.endpoint.deleteAfterRead=false已经满足不删除源文件的要求,要实现无重复消费,补充幂等性相关配置即可:
- 新增配置
camel.source.endpoint.idempotent=true:开启S3连接器内置的幂等机制,自动记录已经读取过的S3对象键,避免重复读取 - 新增配置
camel.source.endpoint.idempotentKey=\${header.CamelAwsS3Key}:指定用S3对象的唯一键作为幂等判断依据 - 配合Kafka Connect原生的偏移量管理机制,新增
offset.flush.interval.ms=60000、errors.tolerance=all配置即可保障消费 Exactly Once 语义。
需求2:仅读取指定文件夹内容
可以通过配置实现,你当前的配置问题是前缀配置错误导致匹配范围不对,修改camel.source.endpoint.prefix配置即可:
- 如果你要读取桶根目录下名为
WriteTopic的文件夹,直接将配置改为camel.source.endpoint.prefix=WriteTopic/,注意结尾的斜杠不可省略,否则会匹配所有名称以WriteTopic开头的文件和文件夹 - 你当前配置的
full/path/to/WriteTopic2会匹配S3桶中路径前缀完全符合该值的对象,如果目标文件夹没有嵌套路径,直接写目标文件夹名加斜杠即可。
修改后的完整配置示例
name=CamelSourceConnector connector.class=org.apache.camel.kafkaconnector.awss3.CamelAwss3SourceConnector key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.camel.kafkaconnector.awss3.converters.S3ObjectConverter camel.source.maxPollDuration=10000 topics=ReadTopic # 配置要读取的S3文件夹前缀,结尾加斜杠精确匹配文件夹 camel.source.endpoint.prefix=WriteTopic/ camel.source.path.bucketNameOrArn=BucketName camel.source.endpoint.autocloseBody=false # 不删除源文件 camel.source.endpoint.deleteAfterRead=false # 开启幂等避免重复消费 camel.source.endpoint.idempotent=true camel.source.endpoint.idempotentKey=\${header.CamelAwsS3Key} # Kafka Connect 偏移量配置 offset.flush.interval.ms=60000 errors.tolerance=all camel.sink.endpoint.region=xxxx camel.component.aws-s3.accessKey=xxxx camel.component.aws-s3.secretKey=xxxx
内容的提问来源于stack exchange,提问作者Adam Gwóźdź
相关产品推荐
相关产品推荐

