
摘要Spark Streaming 的配置参数几十个但真正需要你动手调的就那么几个剩下的调错了反而添乱。这篇把参数按功能分六类重点讲五个高频参数blockInterval、反压、concurrentJobs、unpersist、stopGracefully的默认值、什么时候该动、以及动了之后会踩的坑。关键词Spark Streaming, 配置参数, blockInterval, 反压, concurrentJobs, unpersist一、参数不少但别乱调先给个总判断Spark Streaming 的大部分参数有合理的默认值没搞清楚作用之前别乱动。真正需要你关心的是下面这五个和吞吐、延迟、稳定性直接相关的参数以及它们的配合关系。参数配置有三种方式# ① 提交时作业专属参数用这个spark-submit--confspark.streaming.xxxyyy...# ② 配置文件通用参数放这里spark-defaults.conf# ③ 代码里batchDuration 只能在这里val sscnew StreamingContext(conf, Seconds(2))注意batchDuration在构造函数里设置不能用--conf覆盖这是很多人试了没效果的原因。二、blockIntervalReceiver 模式的分区开关spark.streaming.blockInterval默认 200ms它决定多久生成一个 Block。而 Block 数直接决定分区数分区数 batchDuration / blockInterval比如 batchDuration2s、blockInterval200ms一个 batch 就有 10 个 Block、10 个分区。坑吞吐小的时候如果 blockInterval 设得太小会产生一大堆小 Block调度开销反而比处理开销还大。反过来吞吐大时可以把它调小200ms→100ms让分区更多、并行度更高。判断这是 Receiver 模式的参数。Direct 模式下分区数由 Kafka 分区数决定blockInterval调了没意义。三、反压backpressure.enabled默认false。开启后Spark 会根据当前批处理速度动态调整拉取速率避免数据在内存里越积越多。spark.streaming.backpressure.enabledtruespark.streaming.kafka.maxRatePerPartition10000# 初始上限坑反压不是瞬时生效的它靠 PID 控制器逐步收敛突发流量涌进来的初期仍可能堆积。所以别指望开了反压就一劳永逸还是得配合一个合理的maxRatePerPartition初始值兜底。判断处理能力波动大的场景开反压处理能力稳定的场景手动设maxRatePerPartition就够了。四、concurrentJobs别靠它解决堆积默认1即同一时刻只跑一个 batch 的 job。当前一个 batch 没处理完后一个 batch 的 job 排队等。有些人的第一反应是堆积了就把 concurrentJobs 调大让多个 batch 并行。这通常是个坑多个 job 并行会抢同一批 Executor 资源单个 job 反而更慢输出顺序可能乱有状态场景如 updateStateByKey语义会出问题。判断堆积了正确解法是加 Kafka 分区/加 executor 核数/调大 batchInterval而不是开并发 job。concurrentJobs默认 1 保持不动。五、unpersist内存回收默认true每个 batch 处理完自动 unpersist 掉 RDD释放内存。这是好事绝大多数情况保持默认。坑如果你手动cache了某个 RDD 想跨 batch 复用默认的 unpersist 会把它误清掉导致下个 batch 重算。真有这种跨 batch 缓存的需求才需要把它设成false——但这种情况很少遇到时先想想是不是设计有问题。六、stopGracefullyOnShutdown优雅停机默认false。设成true后应用收到停机信号会先处理完当前 batch 再停不丢数据。spark.streaming.stopGracefullyOnShutdowntrue坑光设这个参数还不够代码里得配合 shutdown hook 才能生效sys.addShutdownHook{ssc.stop(stopSparkContexttrue,stopGracefullytrue)}ssc.start()ssc.awaitTermination()判断生产环境必开。否则每次停机发版、扩容都可能丢当前 batch 正在处理的数据。七、其余几类各挑重点按功能分六类剩下的快速过一遍数据接收maxRatePerPartitionDirect 每分区限流、receiver.maxRateReceiver 每秒上限。状态/checkpointcheckpoint 目录有状态算子 Driver HA 必需、receiver.writeAheadLog.enableReceiver 模式 WALDirect 用不上。容错spark.yarn.maxAppAttemptsYARN 上 Driver 最大重启次数配合--supervise。八、总结参数不少但真正要动的就五个blockInterval、反压、concurrentJobs、unpersist、stopGracefully。blockInterval 决定 Receiver 模式分区数Direct 模式调了没用。反压有收敛延迟得配合 maxRatePerPartition 兜底。堆积了别开 concurrentJobs正确解法是加并行度或调大 batchInterval。unpersist 默认 true跨 batch 缓存才改 falsestopGracefully 生产必开但要配 shutdown hook。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践