如何实时反归一化Aurora PG CDC数据并同步至OpenSearch?
解决流场景下CDC数据反归一化的方案
针对你当前的AWS架构,有几个适配的方案可以实现实时构建反归一化的可搜索文档,替代之前ksqldb的中间件关联逻辑:
方案1:Lambda + DynamoDB 状态存储关联
利用现有Lambda组件,结合DynamoDB保存各表的最新数据快照,实现实时关联:
- 当Lambda收到某张表的CDC事件(插入/更新/删除)时,先将该条数据以主键为Key更新到DynamoDB对应的表条目(比如用
tableName:primaryKey作为全局唯一Key) - 根据反归一化文档的结构,从DynamoDB查询所有需要关联的表的最新数据
- 拼接成完整的反归一化文档后,写入OpenSearch(更新操作直接覆盖原有文档,删除操作同步删除OpenSearch中的对应文档)
- 关键注意点:
- 用DynamoDB的条件更新(如
ConditionExpression)避免并发更新导致的数据不一致 - 为DynamoDB表配置TTL,自动清理过期的历史数据,控制存储成本
- 用DynamoDB的条件更新(如
这个方案的优势是复用现有架构,无需引入新的服务组件,适合关联逻辑相对简单的场景。
方案2:Kinesis Data Analytics (KDA) 流SQL关联
KDA是AWS原生的流处理服务,支持SQL语法,和你之前用的ksqldb逻辑类似:
- 将现有Kinesis Data Stream作为KDA的输入源,为每张表的CDC事件创建对应的流表(定义Schema映射CDC的结构,包括操作类型、主键、字段值等)
- 编写SQL语句实现关联:
- 对于维度表(数据变更频率低),用Lookup JOIN关联流中的事实表事件和维度表的最新快照(KDA会自动维护维度表的状态)
- 对于事实表之间的关联,用窗口JOIN或基于事务ID的分组JOIN,确保同一事务内的变更事件都被关联
- 将JOIN后的反归一化数据输出到新的Kinesis Stream,再通过Lambda写入OpenSearch,或者直接使用KDA的OpenSearch输出连接器
这个方案无需自己管理状态存储,适合纯SQL可实现的关联逻辑,实时性和可靠性都有保障。
方案3:Glue Streaming ETL 自定义关联
如果你的关联逻辑复杂(比如需要自定义转换、多维度嵌套关联),可以用Glue Streaming ETL:
- 配置Glue Streaming Job读取Kinesis Data Stream的CDC数据,定义动态帧(DynamicFrame)解析不同表的事件结构
- 利用Glue的状态管理功能或外部存储(如DynamoDB)维护各表的最新数据状态
- 编写PySpark代码实现自定义关联逻辑,拼接成反归一化文档
- 直接将处理后的结果写入OpenSearch
这个方案灵活性最高,适合大数据量、复杂转换的场景,但需要一定的Spark开发经验。
通用注意事项
- 事件顺序:确保DMS按事务顺序发送CDC事件,Kinesis会保证分片内的事件顺序,处理时尽量按分片顺序处理,避免数据关联错误
- 事务一致性:如果多个表的变更属于同一事务,可以通过CDC事件中的事务ID分组,等待同一事务的所有事件都到达后再生成反归一化文档,保证数据一致性
- 删除事件处理:要同步处理删除事件,避免OpenSearch中残留无效数据,可以通过DynamoDB的状态判断,删除对应文档
内容的提问来源于stack exchange,提问作者Hugo Marcelo Del Negro
相关产品推荐
相关产品推荐

