Spark(Kafka)Streaming内存不足问题排查求助
解决Spark Streaming测试时JVM内存不足的问题
这种情况我在维护Spark 2.x集群做Streaming测试时也碰到过好几次,明明数据量极小却报OOM,十有八九是之前的测试任务没彻底清理残留进程占了内存,给你几个实用的排查和解决步骤:
1. 确认并清理残留的Spark进程
Spark任务如果没有优雅停止,很容易留下SparkSubmit、ExecutorBackend这类后台进程,悄悄占用内存:
- 用
jps命令快速查看当前机器上的Java进程,找到和Spark相关的残留进程后,直接用kill -9 <进程ID>杀掉(注意别误杀其他服务)。 - 更稳妥的方式是用Spark自带的任务管理功能:先通过Spark UI(默认访问
http://<driver-ip>:4040)找到残留任务的application ID,然后执行:
如果是本地测试用的/usr/local/spark/bin/spark-submit --kill <application-id> --master <你的master地址>local[*]模式,master地址填local[*]即可。
2. 调整测试任务的内存配置
Spark默认的内存配置对小数据量测试来说有点过剩,反而容易因为JVM占用过多内存触发系统OOM,提交任务时可以显式限制内存:
/usr/local/spark/bin/spark-submit \ --driver-memory 512m \ --executor-memory 512m \ --class <你的主类> \ <你的jar包路径>
如果是本地模式,Driver和Executor是同一个进程,--driver-memory就能直接限制整个任务的内存使用。
3. 限制Kafka数据源的读取速率
虽然你说测试数据量小,但如果Kafka分区里有积压的旧消息,Spark Streaming可能一次性拉取过多数据导致内存飙升。在Spark 2.2.1的Direct Stream配置里,可以加上:
// Scala示例,可根据你的开发语言调整 val stream = KafkaUtils.createDirectStream(...) stream.foreachRDD(rdd => { // 自定义处理逻辑 }) // 限制每秒从每个分区读取的消息数,测试阶段设小值即可 spark.conf.set("spark.streaming.kafka.maxRatePerPartition", "100")
这样能避免一次性拉取过量数据,从源头减少内存压力。
4. 养成优雅停止任务的习惯
测试时别直接用Ctrl+C强制终止任务,最好通过Spark UI的「Kill」按钮,或者用spark-submit --kill命令停止,让任务有时间清理Kafka的offset、关闭连接、释放内存资源。
内容的提问来源于stack exchange,提问作者TTT
相关产品推荐
相关产品推荐

