spark python3_Python3 连接spark,spark集群 [亲测]
Python3连接Spark集群的关键步骤包括初始化SparkSession、配置Hive支持、设置环境变量及提交作业。单机调试可用简单连接;集群模式需指定Master地址,YARN下通过spark-submit提交,注意PYSPARK_PYTHON等环境变量一致性,并根据需求调整内存、执行器等参数。
先说说怎么回事。Spark和Python的组合,在目前的大数据处理场景下可以说相当常见。对于很多数据工程师或算法工程师来说,掌握了如何在Python环境中灵活地初始化SparkSession,再配合集群提交任务,基本算是日常开发的核心技能了。
1. 连接Spark
连接Spark本身并不是什么特别复杂的事,关键在于细节。从最简单的单机模式到需要支持Hive的集群环境,每一步其实都有值得注意的地方。
1.1. 简单连接
如果你只是想快速在本地验证一下代码逻辑,几行代码就能搞定:
from pyspark.sql import SparkSession
spark = SparkSession \
.builder \
.appName('my_first_app_name') \
.getOrCreate()
这种做法启动快,适合开发初期的调试。
1.2. 连接Spark集群
但如果要正式跑在集群上,特别是需要访问Hive表的话,记得打开Hive支持:
# 使支持hive
spark = SparkSession \
.builder \
.enableHiveSupport() \
.master("xxx.xxx.xxx.xxx:7077") \
.appName("my_first_app_name") \
.getOrCreate()
这里的.master()指向的就是Spark集群的Master地址。需要注意的是,不同的集群管理器(如Standalone、YARN或Mesos)对这里的地址格式要求会有区别。
1.3. 集群配置
在实际生产环境里,比如你的任务需要跑在YARN上,或者是用Standalone模式,常常需要设置SPARK_HOME、PYSPARK_PYTHON这些环境变量。这些配置不光可以在spark-env.sh里写,也可以在Python代码启动前直接指定。
import os
os.environ['SPARK_HOME'] = '/usr/local/workspace/spark-2.1.0-bin-hadoop2.7'
os.environ['PYSPARK_PYTHON'] = '/usr/local/bin/python3.5'
os.environ['PYSPARK_DRIVER_PYTHON'] = 'python3'
from pyspark.sql import SparkSession
spark = SparkSession \
.builder \
.enableHiveSupport() \
.master("xxx.xxx.xxx.xxx:7077") \
.appName("my_first_app_name") \
.getOrCreate()
这里有个小细节:PYSPARK_PYTHON指定的是执行器(executor)上使用的Python,而PYSPARK_DRIVER_PYTHON是指定驱动器(driver)上的。如果你用的是Python虚拟环境,务必保证这两个环境是一致的,不然很容易出现莫名其妙的环境报错。
1.4. config参数
除了这些基础的配置外,在构建SparkSession的时候,还可以通过链式调用.config()来传入各种自定义参数。无论是调整内存大小、设置序列化方式,还是开启一些实验性的特性,都非常方便:
from pyspark.sql import SparkSession
spark = SparkSession \
.builder \
.enableHiveSupport() \
.master("xxx.xxx.xxx.xxx:7077") \
.appName("my_first_app_name") \
.config('spark.some.config.option', 'value') \
.config('spark.some.config.option', 'value') \
... \
.getOrCreate()
这种方式比单独的spark-submit --conf要灵活许多,特别适合在脚本内部根据业务逻辑动态调整配置的场景。
2. 提交作业
提交作业一般有两种方式:一种是我们前面提到的,在代码里直接连接Spark并执行操作;另一种则是通过spark-submit脚本将写好的.py文件提交到集群去执行。
在实际工作中,后者的使用频率其实更高,因为它与集群调度器(如YARN)的集成更自然。
# 提交spark作业
PYSPARK_DRIVER_PYTHON=/opt/anaconda3/envs/xxljob/bin/python \
PYSPARK_PYTHON=/opt/anaconda3/envs/xxljob/bin/python \
/usr/local/workspace/spark-2.1.0-bin-hadoop2.7/bin/spark-submit \
--master yarn \ # 也可以是yarn-client、yarn-cluster
--queue ai \
--num-executors 12 \
--driver-memory 30g \
--executor-cores 4 \
--executor-memory 32G \
/tmp/test_spark.py
需要注意几个关键点:--master yarn这里可以是yarn-client或者yarn-cluster,区别在于Driver进程运行在本地还是集群里。如果是在测试阶段或者交互性要求高,用yarn-client会方便些;如果是生产环境的定时任务,更推荐yarn-cluster,稳定可靠。
此外,如果代码中使用了Python的自定义模块或第三方库,并且环境不是全局安装的,那前面的PYSPARK_PYTHON和PYSPARK_DRIVER_PYTHON就派上大用场了。它可以帮助你在提交任务的同时指定一个隔离的Python虚拟环境,避免与系统Python产生冲突。
Windows 10 是一款微软推出的经典操作系统,拥有硬件兼容性与多任务处理能力。它更偏向把系统状态查看和常用调节动作放在一起,适合需要持续观察和微调设备状态的场景。
极度公式是一款跨平台专业LaTeX公式识别编辑软件,支持OCR公式识别和多平台编辑。和使用说明,避免使用,享受完整功能与稳定支持。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。
















