在 Google Cloud 上运行 Serverless Apache Spark:架构选择、成本调优与 AI 故障排查

2026-08-20 44 预计阅读时间: 1 分钟
来源: cloud.google.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:15 分钟

Apache Spark 仍然是企业级大数据处理的核心工具,但集群创建、YARN 调优、节点扩缩容以及闲置计算资源,会把数据工程团队拖入基础设施运维。Google Cloud 的 Managed Service for Apache Spark 提供了托管集群和 Serverless 两种路径,让团队可以根据工作负载、生态依赖、延迟要求和成本模型做出更准确的选择。

本文围绕三个实际问题展开:什么时候应该使用托管集群,什么时候适合 Serverless;如何控制 Serverless 批处理的性能和 DCU 成本;生产任务失败时,如何借助 Gemini Cloud Assist 缩短定位时间。

先决定部署模型:托管集群还是 Serverless

Serverless 更适合间歇性和事件驱动任务

Serverless 批处理按任务运行期间使用的计算资源计费。任务启动后动态获得执行资源,处理结束后资源释放,因此适合以下场景:

  • 每小时、每天或由工作流临时触发的 ETL 任务。
  • 流量具有明显峰谷的批处理作业。
  • 临时分析、数据回填和一次性数据处理。
  • 团队希望减少集群创建、YARN 配置和节点维护工作。

对于长期运行、资源利用率稳定且接近全天高负载的任务,持续运行的托管 Spark 集群可能更容易控制延迟和成本。如果 Serverless 的启动时间会影响严格 SLA,也需要在选型阶段进行实际测量。

这些情况通常需要托管集群

Serverless 隐藏了底层 VM 和集群管理细节,换来的代价是基础设施控制权减少。以下需求更适合托管集群:

  • 依赖 Flink、Presto/Trino、Hive LLAP 或 HBase 等其他生态组件。
  • 仍然运行 Spark 2.x 代码。
  • 需要自定义 OS 初始化、root SSH、特定本地 SSD 或机器类型。
  • 需要深入调整节点级硬件和操作系统参数。
  • 需要让集群长期保持预热状态,以满足极低延迟要求。

如果只需要在应用层打包特殊库,Serverless 仍然可以通过自定义 Docker 镜像满足一部分依赖需求。但这不等同于获得 VM 层面的控制权。

Serverless 交互式会话与批处理

选定 Serverless 后,还要区分开发方式和生产执行方式。

交互式会话面向人机协作。工程师可以在 Notebook 或开发环境中逐段运行代码,检查 DataFrame、修改变量并生成可视化结果。会话为了保证响应速度,可能在工程师思考期间保持资源活跃,因此长时间闲置会产生不必要的费用。

Serverless 批处理面向自动化运行。它执行完整的 PySpark 脚本或 Java/Scala JAR,适合由 Managed Service for Apache Airflow、Cloud Scheduler、CI/CD 流水线等系统触发。任务结束后计算资源释放,成本边界更清晰。

一个实用的开发流程是:

  1. 使用交互式会话探索数据、确认字段类型并验证转换逻辑。
  2. 将稳定逻辑整理成参数化的 .py.jar 应用。
  3. 使用 Serverless batch 执行生产任务。
  4. 通过编排系统传入输入路径、输出路径和运行日期。

这种方式把“探索时的即时反馈”和“生产时的自动化执行”分开,既方便调试,也避免长期保留开发会话。

用运行时参数控制性能和 DCU 成本

Serverless 并不意味着不需要调优。默认资源规格对简单任务足够,但内存密集型、计算密集型和大规模 Shuffle 作业通常需要显式设置运行时参数。

资源形状要和瓶颈匹配

可以这样判断:

  • 如果任务频繁出现 OOM,优先检查 spark.driver.memoryspark.executor.memory
  • 如果 CPU 长时间满载而内存使用率较低,调整 spark.driver.coresspark.executor.cores
  • 修改 cores 后,最好同时明确设置 memory,避免自动按默认 vCPU 与内存比例分配出不符合预期的资源。

下面是一个可改造的 Serverless Spark 提交示例。运行前请替换项目、区域、代码路径和参数名称;具体提交参数以当前环境中的 Managed Service for Apache Spark CLI 版本为准。

gcloud dataproc batches submit pyspark gs://YOUR_BUCKET/jobs/transactions.py \
  --region=us-central1 \
  --project=YOUR_PROJECT_ID \
  -- \
  --input=gs://YOUR_BUCKET/input/transactions/ \
  --output=gs://YOUR_BUCKET/output/transactions/

脚本内部可以读取这些参数,并在提交时使用 Spark 配置限制资源范围:

import argparse
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, coalesce, lit, try_cast


def parse_args():
    parser = argparse.ArgumentParser()
    parser.add_argument("--input", required=True)
    parser.add_argument("--output", required=True)
    return parser.parse_args()


def main():
    args = parse_args()
    spark = SparkSession.builder.appName("transactions-etl").getOrCreate()

    raw = spark.read.option("header", True).csv(args.input)
    cleaned = (
        raw.withColumn("amount_num", try_cast(col("amount"), "double"))
        .withColumn("quantity_num", try_cast(col("quantity"), "double"))
        .filter(col("amount_num").isNotNull())
        .filter(col("quantity_num").isNotNull())
        .filter(col("quantity_num") != lit(0))
        .withColumn("unit_amount", col("amount_num") / col("quantity_num"))
    )

    cleaned.write.mode("overwrite").parquet(args.output)
    spark.stop()


if __name__ == "__main__":
    main()

上面的 try_cast 用于把异常字符串转换为可处理的空值,随后通过过滤条件跳过无效记录。不同 Spark 运行时对函数支持情况可能不同,部署前应在目标运行时验证;如果不可用,可以改用显式类型转换配合正则校验或异常隔离逻辑。

实际提交时,可以这样配置资源和动态扩缩容边界:

gcloud dataproc batches submit pyspark gs://YOUR_BUCKET/jobs/transactions.py \
  --region=us-central1 \
  --properties=^;^spark.driver.cores=4;spark.driver.memory=16g;spark.executor.cores=4;spark.executor.memory=16g;spark.dynamicAllocation.maxExecutors=20;spark.sql.shuffle.partitions=800 \
  -- \
  --input=gs://YOUR_BUCKET/input/transactions/ \
  --output=gs://YOUR_BUCKET/output/transactions/

命令中的分隔符写法是为了避免多个属性之间的逗号产生歧义。提交前应使用 gcloud dataproc batches submit pyspark --help 检查本地 CLI 对 --properties 的解析方式,并根据实际版本调整格式。

给动态分配设置上限

Serverless 会根据积压任务动态增加 executor,但不受约束的扩容可能放大代码缺陷带来的成本。例如,错误的笛卡尔积、重复重试或异常循环都可能让任务持续申请资源。

spark.dynamicAllocation.maxExecutors 可以作为预算护栏:

  • SLA 优先的任务可以设置较高上限,让任务尽快完成。
  • 夜间批处理可以设置较低上限,接受更长运行时间以换取可预测成本。
  • 生产环境应结合历史运行时间、失败重试策略和输入数据规模设定,而不是直接使用一个过大的固定值。

根据数据规模调整 Shuffle 分区

groupByjoindistinct 等宽转换会产生 Shuffle。默认的 spark.sql.shuffle.partitions 并不一定适合所有数据规模。分区过少时,单个分区可能过大,导致 executor 内存不足并频繁落盘;分区过多则会增加调度和元数据开销。

可以把每个分区约 100 MB 到 200 MB 作为初始估算,再通过 Spark UI 或任务历史数据迭代调整。例如,预计需要处理约 80 GB 的 Shuffle 数据时,可以先从几百到接近一千个分区的范围进行测试,而不是盲目固定为 200。

Google Cloud 提供的基于历史运行数据的 autotuning 会把重复批处理任务按 cohort 归组,并参考过去运行的遥测和统计信息寻找瓶颈。它适合减少人工试错,但仍应保留合理的 executor 上限,并监控自动调优后的运行时间、Shuffle 和费用变化。

用 Gemini Cloud Assist 缩短故障定位路径

生产任务失败时,最浪费时间的往往不是修复代码,而是在 driver、executor 和系统日志之间来回切换。Managed Service for Apache Spark 在 Google Cloud 控制台中集成 Gemini Cloud Assist 后,工程师可以从错误日志直接发起调查,让助手结合日志和任务上下文解释失败原因。

典型的排查过程可以分成三步:

1. 先确认提交参数

如果任务只返回 Application failed with exit code 1,可以从失败日志选择 Investigate log,让助手检查 driver 日志和执行上下文。常见原因是提交时遗漏了输入 GCS 路径、输出路径或日期参数,而脚本使用 argparse 将这些参数声明为必填项。

2. 再检查 schema 和脏数据

补齐参数后,任务可能进入下一层失败:CSV 自动推断把 amountquantity 当成字符串,而代码却直接进行除法。源文件中还可能混有空字符串、文本金额或零值分母。

排查时应同时确认两件事:转换表达式的输入类型,以及原始文件中是否存在不符合预期的记录。只修改 DataFrame 表达式而不检查输入数据,往往会让同类问题在下一批文件中再次出现。

3. 让助手生成修复建议,但保留人工验证

可以向 Gemini Cloud Assist 提出类似请求:

请将计算逻辑改为 amount / quantity,为两个字段增加安全类型转换;跳过无法转换、为空或 quantity 为零的记录,同时说明被跳过记录的数量应该如何监控。

助手可以生成类型转换、空值处理和过滤逻辑,但上线前仍需要工程师验证:

  • 业务语义是否确实是 amount / quantity
  • 丢弃无效记录是否符合数据质量策略。
  • 是否需要把坏记录写入隔离路径,而不是直接丢弃。
  • 修复后结果数量、金额汇总和 SLA 是否符合预期。

AI 适合缩短日志阅读和修复草稿生成时间,不能替代数据契约、回归测试和发布审批。

落地时的检查清单

采用 Serverless Apache Spark 前,可以按下面的顺序评估:

  • 工作负载是否间歇性、突发性或由编排系统触发?
  • 是否依赖 Spark 3.x+,且不需要 VM、OS 或 root 级别定制?
  • 开发阶段是否使用交互式会话,生产阶段是否切换到参数化批处理?
  • 是否明确设置 driver/executor 的 cores 和 memory?
  • 是否为动态扩缩容设置 maxExecutors
  • 是否根据 Shuffle 数据量调整 spark.sql.shuffle.partitions
  • 是否记录每次运行的耗时、输入规模、Shuffle、失败原因和 DCU 成本?
  • 是否为 schema 异常、空值、零值和坏记录设计处理策略?
  • 是否把 Gemini Cloud Assist 的建议纳入人工验证和测试流程?

Serverless 的价值不只是“无需管理集群”,而是把基础设施决策从日常运维中抽离出来,同时通过运行时参数、扩缩容上限和历史调优保留成本控制能力。对于突发批处理和自动化 ETL,它可以减少闲置资源;对于需要深度生态兼容或底层控制的长期任务,托管集群仍然是更稳妥的选择。


相关推荐