You need to enable JavaScript to run this app.
文档中心
E-MapReduce

E-MapReduce

复制全文
下载 pdf
Serverless Spark 作业:基于 Spark connect
添加 Spark Connect App 依赖
复制全文
下载 pdf
添加 Spark Connect App 依赖

背景信息

使用 Spark Connect 执行任务时,代码实际由远端 Spark Connect App 的 Executor 运行。如果任务中用到了本地自定义库、辅助文件或压缩包,需要将这些依赖显式传入 App,Executor 才能在运行时正确加载。
当前支持两种方式传入依赖:

  • 在控制台创建 App 时通过参数预先配置。
  • 在运行时通过 addArtifact 接口动态上传(仅 Python)。

依赖包较大时建议优先在创建 App 时配置,避免每次运行重复上传。

说明

目前不支持从客户端上传依赖到服务端。

控制台创建 Spark Connect App 时配置依赖

创建指导请参见EMR 控制台创建,并通过参数 Archive 指定依赖包路径,上传依赖。
Image

PySpark 通过 addArtifact 动态上传依赖

Spark Connect 客户端支持通过 addArtifact 接口将 Python 环境和依赖文件上传到 Spark Connect App,再通过 SparkContext 分发到各个 Executor 上。

Python 环境(推荐 Conda 打包)

注意点:

  1. 需要使用linux环境打包
  2. Driver 和 Executor 环境的 Python 版本需要一致,Serverless Spark 默认使用 Python 3.12
  3. 如果 Conda 包比较大,上传时间会很长,建议在创建 Spark Connect App 时配置
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")

File

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)

Pyfile(支持 Python 文件和 zip 压缩包)

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)

Archive

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)

最近更新时间:2026.08.14 15:46:15
这个页面对您有帮助吗?
有用
有用
无用
无用