本节摘要:DStream 的算子分为无状态转换(每个批次独立计算)与输出操作(触发小批作业的引擎钩子)。本节按执行视角过一遍高频算子的批次展开语义,用转换链检查工具观察模板,并梳理输出操作与检查点目录的配合。
无状态算子只看当前批次,批次之间零记忆。写法与 RDD 完全同构——因为它们本来就是逐批次套用的 RDD 算子:
from pyspark import SparkContext from pyspark.streaming import StreamingContext sc = SparkContext(master="local[2]", appName="DstreamOps") ssc = StreamingContext(sc, 3) events = ssc.socketTextStream("localhost", 9999) # 管道组合:解析 -> 过滤 -> 计数,全部只作用于当前批次 per_batch = (events.map(lambda line: line.split(",")) .filter(lambda f: f[1] == "click") .map(lambda f: (f[0], 1)) .reduceByKey(lambda a, b: a + b)) # 引擎视角:每个批次间隔,对当批 RDD 依次施加 map filter map reduceByKey per_batch.pprint()
值得巡检的一点:reduceByKey 在每个小批作业里都触发一次 Shuffle。批次间隔越短、Shuffle 越频繁——流作业的 Shuffle 开销是"每秒都在付"的,比批作业敏感得多。能用 reduceByKey 就不用 groupByKey 的批作业守则,在流侧更要严格。
transform 算子给你开后门的权利:拿到每个批次的裸 RDD,混用 RDD 级操作甚至引用外部广播变量:
def enrich(rdd): return rdd.map(lambda kv: (kv[0], kv[1] * 10)).filter(lambda kv: kv[1] > 50) enriched = per_batch.transform(enrich) # 每批次把函数套在 RDD 上
转换模板静静躺着,直到输出操作出现——它才是"行动算子"的流侧对应物,把当批 RDD 提交为真正的小批作业:
| 输出操作 | 触发时机 | 典型用途 |
|---|---|---|
| pprint | 每批次 | 调试观察 |
| foreachRDD | 每批次 | 写外部系统、任意副作用 |
| saveAsTextFiles | 每批次 | 按时间戳落盘归档 |
| count / collect 系 | 每批次 | 取回 Driver 端小结果 |
foreachRDD 是生产主力,也是事故高发区:
def send(batch_rdd): if not batch_rdd.isEmpty(): # 连接必须在 Task 内创建:连接对象不可序列化,也不该共享 batch_rdd.foreachPartition( lambda part: _write_to_sink(list(part))) events.foreachRDD(send)
巡检守则有两条。第一,连接在分区内创建复用,不在 Driver 建、不在逐条建——前者传不进 Task,后者每条记录一次握手。第二,写外部系统时要考虑重复:微批重跑时 foreachRDD 会再执行一遍,幂等写入或事务批次号是配套必修课(3.4 节展开)。
print(per_batch.dstream.) # 更直观的方式:在 UI 上比较相邻两个小批作业的 Stage 拓扑——完全相同 # 因为它们来自同一份模板,只是数据换了批次
UI 上相邻作业的 DAG 拓扑逐字相同,是"模板展开"最直接的证据。若某个作业突然多出一个 Stage,先怀疑是 transform 里的条件逻辑在不同批次走了不同分支。
💡 关键直觉:DStream 代码是"批处理程序的图纸",引擎每隔一个批次间隔照图纸盖一栋楼。调试流作业的许多问题,抽出一个批次单独当批作业复现,比盯着流动的日志高效得多。
背景:运营大屏每 3 秒刷新各渠道点击量,管道为 socket 进来的事件做解析、过滤、按渠道计数后经 foreachRDD 写入 Redis。上线一周后偶发"大屏卡 40 秒不动"。操作:先固定证据——把 UI 切到作业列表按时间排开,发现卡顿时段的小批作业并未失败,而是排队:某几个批次的作业耗时冲到 20 秒以上,后续批次在接收器侧积压。定位:批内数据不均,整点活动开始时点击洪峰把单批数据量抬了 8 倍,reduceByKey 的 Shuffle 与 Redis 写入同步变慢。处置分两层:引擎侧把批次间隔从 3 秒调到 5 秒并给 receiver 的最大速率限流,让单批输入封顶;写入侧把逐条写 Redis 改为分区内 pipeline 批量写,单批写耗时从 14 秒降到 2 秒。结果:洪峰期端到端延迟稳定在 8 秒内,大屏恢复流动。解读:微批系统的"卡住"通常不是死,是积压排队——先看作业列表的时间线分布再谈代码,一眼区分"作业慢"与"作业排队"。变式:若洪峰不可预测且延迟要求苛刻,调大间隔只是止血,根治要换反压机制更强的接收结构或更细的分区键。
# 写入侧的改造核心:分区级 pipeline 批量写 def send(batch_rdd): if batch_rdd.isEmpty(): return def write_part(part): import redis r = redis.Redis(host="sink", port=6379) # 连接在分区内创建 pipe = r.pipeline(transaction=False) for channel, cnt in part: pipe.hincrby("rt:click", channel, cnt) # 幂等增量 pipe.execute() # 一次网络往返 batch_rdd.foreachPartition(write_part)