如何基于Debezium消息操作类型配置S3 Sink Connector路由规则
Debezium Kafka消息按操作类型路由至S3分桶配置方案
路由规则映射
Debezium输出的变更事件中,操作类型标记存储在消息体的op字段,和目标分桶的对应关系如下:
op = r:全量快照阶段生成的读取事件,写入bucket 1op = u(更新事件)、op = d(删除事件):增量同步阶段生成的变更事件,写入bucket 2
该能力不需要额外部署流处理组件做中转,直接通过Kafka Connect S3 Sink Connector内置的单消息转换(SMT)能力即可实现。
实现方案
优先选择双Sink实例+消息过滤的实现方式,配置简单、故障隔离性强,排障成本低。
1. 快照数据Sink配置(写入bucket 1)
核心配置片段,其余S3认证、写入格式、分区规则按常规S3 Sink要求补全即可:
name=debezium-s3-snapshot-sink connector.class=io.confluent.connect.s3.S3SinkConnector s3.bucket.name=<替换为实际bucket1名称> topics=<替换为Debezium输出的源Topic名称> # 配置消息过滤逻辑,仅保留op=r的快照事件 transforms=filterNonSnapshot transforms.filterNonSnapshot.type=org.apache.kafka.connect.transforms.Filter transforms.filterNonSnapshot.predicate=matchReadOp # 定义谓词匹配规则:匹配op字段值为r的消息 predicates=matchReadOp predicates.matchReadOp.type=org.apache.kafka.connect.transforms.predicates.RecordValue predicates.matchReadOp.field=op predicates.matchReadOp.pattern=r
2. 增量数据Sink配置(写入bucket 2)
核心配置片段,其余通用配置和快照Sink保持一致即可:
name=debezium-s3-incremental-sink connector.class=io.confluent.connect.s3.S3SinkConnector s3.bucket.name=<替换为实际bucket2名称> topics=<替换为Debezium输出的源Topic名称> # 配置消息过滤逻辑,仅保留op=u、op=d的增量变更事件 transforms=filterNonIncremental transforms.filterNonIncremental.type=org.apache.kafka.connect.transforms.Filter transforms.filterNonIncremental.predicate=matchWriteOp # 定义谓词匹配规则:匹配op字段值为u或d的消息 predicates=matchWriteOp predicates.matchWriteOp.type=org.apache.kafka.connect.transforms.predicates.RecordValue predicates.matchWriteOp.field=op predicates.matchWriteOp.pattern=u|d
如果不想维护两个Sink实例,也可以使用单实例+动态路由SMT的方案,但该方案要求S3 Sink Connector版本支持从消息上下文动态读取bucket配置变量,生产环境不做优先推荐。
内容的提问来源于stack exchange,提问作者Sanjeev Sharma
相关产品推荐
相关产品推荐

