流处理中批处理有何作用?相比逐元素处理有何优势?
流处理场景下微批处理相比逐元素单条处理的优势
相比每条数据到达就立刻处理的模式,把数据流切分成固定大小/固定时间间隔的小批次做批量计算,核心优势有这几点:
- 系统固定开销被大幅摊薄:不管是单条还是一批数据,处理流程里的任务调度、序列化反序列化、网络IO、状态存储读写、容错日志记录这些操作的固定成本基本持平。举个实际场景的例子,单条1KB的数据走一次完整分布式处理流程,固定开销可能占总耗时的90%以上;如果攒成1秒间隔的批次,单批包含几万条同类型数据,同样的固定开销平摊到所有数据上,单条数据的平均处理成本能压到原来的几十分之一,整体吞吐量能拉开一个甚至两个量级的差距。
- 容错实现简单、故障恢复快:如果做逐条处理,要保证精确一次的处理语义,就得给每条数据单独记录处理位点、打校验点,光记录这些日志产生的IO压力就可能把存储集群压垮。微批模式下只需要记录每个批次的起始、结束偏移量,给整个批次做一次检查点即可,出现故障时直接重跑失败的小批次就行,逻辑简单,恢复速度也快很多。
- 计算算子执行效率更高:聚合、关联、排序、特征计算这类流处理常见算子,批量执行时引擎可以做大量底层优化,比如适配CPU缓存的向量化执行、map端预聚合减少跨节点shuffle的数据量、批量压缩降低存储和传输成本等等,这些优化在逐条处理模式下根本无法落地——毕竟单条数据连向量化执行要求的最小批量长度都凑不够。
- 集群资源利用率更平稳:真实业务场景里的数据流到达速率从来不是稳定值,高峰和低谷的流量差可能达到几十上百倍。逐条处理很容易在流量高峰时把计算资源打满导致队列堵塞,流量低谷时又让资源大量空转浪费。按固定间隔攒批的模式可以天然平滑流量波动,让集群的CPU、内存、网络负载保持在稳定区间,避免频繁的资源扩缩容抖动。
为什么无法做到每条数据到达就立刻执行单条处理
本质上不是技术上完全做不到收一条处理一条,而是这种模式在生产环境下的投入产出比极低,甚至根本跑不通实际业务:
- 绝对零延迟本身就不存在:数据从业务端产生,经过采集、传输、消息队列缓存到达计算引擎的链路本身就存在固有延迟,就算引擎收到一条数据就立刻启动计算,也消不掉上游链路的延迟。为了节省引擎侧攒批的几百毫秒延迟,把整体吞吐量打个一折,在99%的业务场景里都是完全不划算的选择。
- 分布式系统的协调成本扛不住:生产环境的流处理基本都是分布式部署,数据要按规则分片、跨节点做数据交换、同步全局状态、做一致性校验,这些协调操作不可能跟着每条数据执行。举个最简单的例子,要是做全局访问量统计时每条数据来都抢全局锁更新计数,整个系统会直接被锁竞争卡死,根本没法正常运行。
- 大部分业务逻辑本身不支持单条独立计算:流处理的常见需求比如时间窗口统计、异常流量检测、规则匹配、流表关联,本身就需要结合一个时间范围内的数据集、或者关联外部存储的批量数据才能算出有效结果,单条孤立的数据根本没法产出有业务意义的计算结果。
Spark Streaming 中的核心抽象
DStream就是微批设计思路的典型实现:它把连续的数据流按用户配置的批间隔(比如1秒、5秒)切分成一个个离散的小批次数据,每个批次直接复用Spark核心的RDD批处理逻辑执行计算,这套设计让Spark Streaming在诞生初期就实现了远超当时逐条处理类流引擎的吞吐量表现,相关设计细节可以查阅Spark官方Streaming编程指南的DStream章节。
内容的提问来源于stack exchange,提问作者Vighnesh Nayak
相关产品推荐
相关产品推荐

