知识卡片
流处理的恰好一次语义:微批、存档点、原子提交与幂等性
内容
[[批处理输出的三种目的与不可变输出带来的人类容错]]里MapReduce能轻松做到恰好一次语义(更贴切的说法是”等效一次”),靠的是等任务全部成功才让输出可见、失败任务的输出直接丢弃——但流处理无法”等它跑完”,因为流永不结束。解决办法之一是微批次(如Spark Streaming):把流切成约1秒的小块当微型批处理处理,批次越小调度协调开销越大、批次越大结果可见延迟越长;Flink则用定期存档点:算子状态定期写入持久存储,一旦崩溃就从最近存档点重启、丢弃存档点之后的输出。但这两种方法只能保证流处理框架内部的恰好一次,一旦输出离开框架(写外部数据库、发消息、发邮件),框架就无法再撤销失败批次已经外泄的副作用了——要让外部副作用也做到恰好一次,需要让处理的全部输出与副作用”当且仅当处理成功时才生效”,这本质上就是[[原子提交问题与两阶段提交的机制及两个不归路]]的思路在流处理场景里的重现(Google Cloud Dataflow、VoltDB已经在更受限的环境里高效实现了这种原子提交,不像XA那样要跨异构技术)。另一条更轻量的路径是依赖幂等性:多次执行和执行一次效果相同的操作(如”把某键设为某值”天然幂等,”计数器加一”则不是,但配合消息的持久递增偏移量作为去重依据,也能把非幂等操作改造成幂等),前提是失败重启必须按相同顺序重放相同消息、处理必须确定性、且不能有其他节点同时更新同一个值,必要时还要用防护机制阻止假死节点干扰。
参考来源
- 位置:《数据密集型应用系统设计》第十一章《流处理》"容错""微批量与存档点""原子提交再现""幂等性"(源文件:_epub-src/ch11_split_002.html)
- 结论依据:原文说明微批次与存档点为流处理框架内部提供恰好一次语义、但无法撤销已外泄的副作用,进而介绍原子提交与幂等性两条实现外部恰好一次效果的路径及各自前提条件,直接支撑本卡片结论。
- 原始内容:一个解决方案是将流分解成小块,并像微型批处理一样处理每个块……只要输出离开流处理器,框架就无法抛弃失败批次的输出了……幂等操作是多次重复执行与单次执行效果相同的操作。