使用 Spark Connect 执行任务时,代码实际由远端 Spark Connect App 的 Executor 运行。如果任务中用到了本地自定义库、辅助文件或压缩包,需要将这些依赖显式传入 App,Executor 才能在运行时正确加载。
当前支持两种方式传入依赖:
依赖包较大时建议优先在创建 App 时配置,避免每次运行重复上传。
说明
目前不支持从客户端上传依赖到服务端。
创建指导请参见EMR 控制台创建,并通过参数 Archive 指定依赖包路径,上传依赖。
Spark Connect 客户端支持通过 addArtifact 接口将 Python 环境和依赖文件上传到 Spark Connect App,再通过 SparkContext 分发到各个 Executor 上。
注意点:
import os from pathlib import Path import conda_pack from pyspark.sql import SparkSession, functions as F from pyspark.sql.types import StringType # 创建 SparkSession spark = SparkSession.builder.getOrCreate() @F.udf(StringType()) def probe(_: int) -> str: import sys return f"python={sys.executable}; version={sys.version.split()[0]}; prefix={sys.prefix}" def show_probe(label: str) -> None: spark.range(1).select(probe(F.lit(1)).alias(label)).show(truncate=False) show_probe("before_addartifact") # 将当前环境的pyspark_conda_env打包成pyspark_conda_env.tar.gz, 也可以通过conda pack命令打包、 # 注意点 # 1. 需要使用linux环境打包 # 2. Driver和Executor环境的python版本需要一致 conda_pack.pack() spark.addArtifact( f"{os.environ.get('CONDA_DEFAULT_ENV')}.tar.gz#environment", archive=True) # 设置executor python路径 spark.conf.set( "spark.sql.execution.pyspark.python", "environment/bin/python") show_probe("after_addartifact")
from pathlib import Path from pyspark.sql import SparkSession, functions as F from pyspark.sql.types import StringType root = Path.cwd() spark.addArtifact(str(root / "hello_connect.txt"), file=True) @F.udf(StringType()) def read_file(_: int) -> str: from pyspark import SparkFiles with open(SparkFiles.get("hello_connect.txt"), "r", encoding="utf-8") as f: return f.read().strip() spark.range(1).select(read_file(F.lit(1)).alias("file_result")).show(truncate=False)
from pathlib import Path from pyspark.sql import SparkSession, functions as F from pyspark.sql.types import StringType root = Path.cwd() spark.addArtifact(str(root / "connect_demo_helper.py"), pyfile=True) # 支持自动解压zip压缩 spark.addArtifact(str(root / "demo_archive.zip"), pyfile=True) @F.udf(StringType()) def use_pyfile(_: int) -> str: import connect_demo_helper return connect_demo_helper.describe_pyfile_artifact() spark.range(1).select(use_pyfile(F.lit(1)).alias("pyfile_result")).show(truncate=False)
from pathlib import Path from pyspark.sql import SparkSession, functions as F from pyspark.sql.types import StringType root = Path.cwd() spark.addArtifact(str(root / "demo_archive.zip"), archive=True) @F.udf(StringType()) def read_archive(_: int) -> str: from pyspark import SparkFiles archive_dir = SparkFiles.get("demo_archive.zip") archive_file = os.path.join(archive_dir, "archive_message.txt") with open(archive_file, "r", encoding="utf-8") as f: return f.read().strip() spark.range(1).select(read_archive(F.lit(1)).alias("archive_result")).show(truncate=False)