如何将Kafka中Debezium消息同步至数据湖与Redshift并确保数据湖为可信源
问题解答
核心问题:Kafka是否自带同时写入Redshift和S3的连接器?
没有。Kafka Connect官方提供的是独立的Sink连接器:
S3 Sink Connector:负责将Kafka消息写入S3数据湖JDBC Sink Connector:通过JDBC协议将消息写入Redshift(需配合Debezium SMT转换格式)
这两个连接器无法合并为一个完成两项任务。
关于双连接器一致性的顾虑
你提到的「同时用两个连接器从同一Debezium Topic读取」的方案,确实存在数据一致性风险:两个连接器各自维护消费偏移量,一旦其中一个出现故障、重试或延迟,会导致S3数据湖和Redshift的数据不一致,无法保证数据湖作为可信数据源的唯一性。
推荐架构实现
你考虑的多Topic链路是更可靠的方案,具体落地可以参考以下步骤:
- 原始消息落地数据湖:使用
S3 Sink Connector将Debezium Topic的所有原始消息完整写入S3,确保数据湖作为可信的单一数据源。 - 转换并生成新Topic:从S3读取原始Debezium消息,提取
after字段转换为行式格式后写入新的Kafka Topic。这一步可以通过以下方式实现:- 用
S3 Source Connector读取S3文件,配合Debezium SMT(如ExtractNewRecordState)转换格式后写入新Topic - 用Flink/Spark流处理工具直接读取S3上的Debezium数据,转换后写入新Topic
- 用AWS Lambda触发S3文件变更事件,转换消息后写入新Topic
- 用
- 写入Redshift:使用
JDBC Sink Connector从新的行式数据Topic读取内容,直接写入Redshift(无需额外转换,因为已经在前置步骤处理完成)
这种架构确保所有数据都先流经S3数据湖,Redshift的数据完全来源于数据湖处理后的结果,从根源保证了数据湖的可信性和链路一致性。
内容的提问来源于stack exchange,提问作者Jwan622
相关产品推荐
相关产品推荐

