DynamoDB Stream数据消费方案选型及shard过期故障应对咨询
DynamoDB 数据同步方案选型分析
核心结论
两种方案没有绝对的好坏,完全匹配你的场景就是最优解,其中你最担心的DynamoDB Stream 24小时保留期丢数据的问题,两种方案都有对应的规避手段,核心差异在链路复杂度和数据可靠性冗余度。
两种方案优劣势对比
方案1:直接用Python boto客户端消费DynamoDB Stream写入自有存储
- 优势:链路最短、额外组件最少,运维成本和资源开销极低,数据同步延迟通常在秒级,适合数据量不大、下游自有存储可用性高的非核心业务场景。
- 劣势:无缓冲层,一旦下游存储故障、或者消费链路卡住超过24小时,还未消费的Stream数据会永久丢失,数据可靠性完全依赖消费链路的稳定性和故障恢复速度。
方案2:新增Kafka/SNS作为中间缓冲层,再从中间层消费写入存储
- 优势:完美解决24小时保留期的限制,Kafka支持自定义消息保留时长,默认可以配置7天甚至更长,就算下游存储故障数天,只要中间层消息没过期都可以回溯消费,完全避免Stream过期丢数据的问题;同时可以做削峰填谷,扛住DynamoDB的突增写入流量,不会直接打垮下游存储;如果后续有其他业务需要消费DynamoDB变更数据,直接接中间层即可,不需要重复消费Stream。
- 劣势:需要额外运维中间件组件,链路复杂度更高,有一定的资源和运维成本,同步延迟比方案1略高。
故障防丢数据的配套措施
不管选哪种方案,都需要做以下配置避免数据丢失:
提前开启DynamoDB Stream全类型变更捕获,设置事件投递为至少一次模式,避免上游变更事件漏发。
选择方案1的额外配置
- 持久化消费进度:用独立的存储(比如Redis、小型DynamoDB表)存储每个分片的最新消费
SequenceNumber,故障恢复后直接从上次消费的位置继续读取,避免浪费时间从头遍历分片。 - 配置死信队列:遇到下游写入失败的异常数据,直接推送到死信队列暂存,不要阻塞整个分片的消费,避免个别异常数据卡住整条链路,最终导致Stream过期丢数据。
- 配置消费积压告警:实时监控每个分片的最新消息序列号和当前消费序列号的差值,超过阈值立刻告警,预留足够的处理时间避免积压超过24小时。
选择方案2的额外配置
- 写入中间件时开启强确认:Kafka配置
acks=all,确保消息成功写入多副本后再提交Stream的消费进度,避免中间层本身丢消息。 - 中间层消息保留时长建议设置不小于7天,同时开启中间层的数据多副本备份,避免中间层故障导致数据丢失。
- 下游消费进度直接存在Kafka自带的
__consumer_offsets主题即可,故障恢复后直接回溯消费即可。
选型建议
如果你的同步数据是非核心数据、数据量小、下游自有存储可用性高,故障恢复时间基本都在12小时以内,选方案1就足够,性价比最高。
如果是核心业务数据,完全不能接受数据丢失,或者下游存储故障恢复时间可能超过24小时,或者数据量很大有明显的写入峰值,优先选择方案2,虽然多了一点运维成本,但数据可靠性冗余度高很多,长期来看更省心。
内容的提问来源于stack exchange,提问作者Avenger
相关产品推荐
相关产品推荐

