十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SparkStreaming 之 DStream 底层结构剖析

SparkStreaming 之 DStream 底层结构剖析 摘要DStream 是 Spark Streaming 的核心抽象但很多人用了很久也没搞清楚它到底是什么。这篇把 DStream 拆开——它本质是RDD 的时间序列本身不存数据讲清它内部的 DStreamGraph、Lineage 血缘以及 getOrCompute 怎么按时间点把 DStream 变成 RDD顺带说清有状态操作为什么要强制开 checkpoint。关键词Spark Streaming, DStream, RDD 序列, DStreamGraph, getOrCompute, checkpoint一、DStream 到底是个什么东西先给一个判断省得绕DStream 是RDD 的时间序列它自己一个字节的数据都不存。回想上一篇文章里 batchInterval 的概念——Spark Streaming 每 2 秒举例切一个 batch每个 batch 就是一个 RDD。DStream 就是把这些按时间排好的 RDD 串起来的抽象RDD t0 RDD t1 RDD t2 ... RDD tn (batch 0) (batch 1) (batch 2) (batch n)所以你写的ds.map(...)、ds.reduceByKey(...)这些操作本质上是在描述每个时间片的 RDD 该怎么变换而不是立刻去算。这一点和 RDD 的惰性求值是一脉相承的——DStream 的转换操作构建的是依赖链只有遇到 output 操作print、saveAsTextFiles、foreachRDD才真正触发计算。二、DStream 的内部结构一个 DStream 内部有三块关键的东西。2.1 DStreamGraphDStream 的 DAG和 RDD 有 DAG 一样DStream 也有自己的图结构DStreamGraph里面按角色分三类节点InputDStream数据入口比如SocketInputDStream、KafkaInputDStream对应上一篇讲的 Receiver。TransformedDStream中间转换map/flatMap/filter/join/window都会产生这种节点。ForEachDStream输出节点foreachRDD、print属于这一类。DStreamGraph维护这张依赖图每个 batch 时间点到了就顺着图把每个 DStream 对应的 RDD 算出来。2.2 Lineage血缘DStream 同样记录完整的血缘——“我是从哪个父 DStream、经过什么转换来的”。这条链子的用途和 RDD 的 lineage 一样容错。不过 DStream 的血缘有个 RDD 没有的要求涉及有状态操作或需要故障恢复时必须开 checkpoint。因为 DStream 的血缘是逻辑上的转换链跨很多 batch 之后这条链会非常长光靠血缘重算既不现实也没法恢复累计到现在的状态。checkpoint 会把依赖链和状态持久化到可靠存储HDFSDriver 挂了能从 checkpoint 拉起来接着算。ssc.checkpoint(hdfs://namenode:8020/checkpoint/streaming-app)// 有状态操作不开 checkpoint 会直接抛异常valstatelines.updateStateByKey(updateFunc)三、转换操作的两个分类维度理解 DStream 的操作可以按两个维度切这直接影响性能和要不要 checkpoint。维度一无状态 vs 有状态无状态有状态代表算子map / filter / union / joinupdateStateByKey / window是否跨 batch只看当前 batch跨 batch 累计状态checkpoint不需要必须开无状态转换每个 RDD 独立处理干净利落。有状态转换要在 batch 之间记住状态这个状态就是靠 checkpoint 落盘保住的。维度二无 shuffle vs 有 shufflemap/flatMap/filter是窄依赖一个分区内独立算不碰网络join/groupByKey/reduceByKey是宽依赖要跨分区 shuffle代价高一个量级。这和 RDD 的判断完全一致——DStream 的每个操作最终都会落到对应 RDD 的某个操作上。四、getOrComputeDStream 怎么变成 RDD这是 DStream 最核心的一个方法把按时间生成 RDD这件事讲清楚。JobGenerator 每个 batch 触发一次传入时间点t然后对每个 DStream 调用getOrCompute(t)查缓存先看generatedRDDs里有没有t对应的 RDD有就直接返回。没命中就 compute(t)调用当前 DStream 的compute它会先递归地让父 DStream 生成t时刻的 RDD再应用当前这一步的转换。缓存返回算出来的 RDD 存进generatedRDDs避免重复计算。为什么要缓存因为同一个 DStream 在同一个时间点可能被多个下游引用——比如一个流既要做实时统计又要写一份原始数据。不缓存的话同一个 RDD 会被算两遍浪费整条链路的计算。五、transform 和 foreachRDD打通 RDD APIDStream 的算子就那么几十个遇到复杂逻辑或者要访问外部资源连接池、状态存储时就不够用了。transform和foreachRDD就是为这个留的口子// transform拿到 RDD返回新的 RDD继续流式计算valcleanedds.transform{rddrdd.mapPartitions{iter/* 每分区初始化一次连接 */}}// foreachRDD拿到 RDD做输出到此为止不再返回ds.foreachRDD{rddrdd.foreachPartition{partitionOfRecords/* 写外部存储 */}}两者都在 Driver 端直接拿到 RDD可以随意用 RDD API。一个常见的最佳实践是把连接这类重量级对象的创建放进mapPartitions/foreachPartition里而不是在每条记录上 new 一个连接否则吞吐会崩。六、总结DStream 是 RDD 的时间序列不存数据转换操作是惰性的。DStreamGraph 分 Input/Transformed/ForEach 三类节点Lineage 负责容错有状态操作必须 checkpoint。转换按无状态/有状态和无 shuffle/有 shuffle两个维度判断前者决定要不要 checkpoint后者决定代价高低。getOrCompute 是 DStream 到 RDD 的关键查缓存 → 递归 compute → 缓存返回缓存是为了不让同一时间点的 RDD 被重复计算。复杂逻辑用 transform/foreachRDD 打通 RDD API连接等重对象放分区级初始化。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表