在Ubuntu虚拟机上使用spark-submit遇错误,求教正确用法
Spark流处理:Jupyter Notebook运行与spark-submit提交指南
一、Jupyter Notebook中正确运行Spark流处理
Spark流处理(包括Structured Streaming)完全可以在Jupyter Notebook中运行,核心是确保Spark环境配置正确:
- 初始化正确的SparkSession
from pyspark.sql import SparkSession # 本地开发模式,生产环境替换为集群master地址(如spark://host:7077) spark = SparkSession.builder \ .appName("LocalStreamingDemo") \ .master("local[*]") \ .getOrCreate()
- 运行流处理逻辑(以Socket流为例)
# 监听本地9999端口的Socket数据 stream_df = spark.readStream \ .format("socket") \ .option("host", "localhost") \ .option("port", 9999) \ .load() # 输出到控制台,启动流查询 query = stream_df.writeStream \ .outputMode("append") \ .format("console") \ .start() # 保持作业运行,直到手动停止 query.awaitTermination()
注意:运行单元格后,Notebook会处于持续运行状态,需手动中断内核停止作业;若使用Kafka等外部数据源,需在SparkSession初始化时添加依赖配置,比如.config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0")
二、spark-submit提交流处理脚本的正确步骤
如果需要将流处理作业提交到集群运行,使用spark-submit的标准流程如下:
编写独立的流处理脚本(如
streaming_job.py),核心逻辑与Notebook中一致执行提交命令
# 本地模式提交 spark-submit \ --master local[*] \ --name ClusterStreamingJob \ streaming_job.py # YARN集群模式提交 spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 2G \ --num-executors 3 \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0 \ streaming_job.py
关键参数说明
--master:指定运行模式,local[*]为本地模式,yarn/spark://host:7077为集群模式--deploy-mode:集群模式下可选cluster(作业在集群节点运行)或client(本地提交节点作为驱动)--packages:加载第三方依赖(如Kafka、Redis连接器),需保证Spark版本与依赖版本兼容--executor-memory/--num-executors:配置集群资源,根据作业需求调整
常见问题排查
- Notebook报错时,优先检查SparkSession是否加载了所需依赖,以及本地端口/数据源是否可访问
- 集群提交失败时,确认Spark集群配置正常,依赖包可通过集群节点获取(或通过
--packages自动下载)
内容的提问来源于stack exchange,提问作者Elmauro
相关产品推荐
相关产品推荐

