Google Cloud 无服务器 Apache Spark:架构与 AI 故障排查

Google Cloud 的 Managed Service for Apache Spark 提供托管集群与 serverless 两种模式。选型依据有三:负载特征、生态系统限制、底层控制需求。持续高频、基线利用率高的 24/7 流水线,传统集群更易预测成本;间歇性、突发性或由编排器触发的任务,serverless 不再为闲置计算付费。如果依赖 Flink、Trino、HBase 或旧 Spark 2.x,则必须选择托管集群;需要 SSH、自定义 OS 初始化或特定存储配置时也一样。serverless 仅支持 Spark 3.x+,但可打包自定义 Docker 镜像应对应用级库需求。

在 serverless 内还要区分交互式会话与批量任务。交互式适合 Notebook 中逐 cell 调试,数据驻留内存,思考期间有闲置成本;batch 运行打包好的 .py 或 .jar,由 Airflow 或 CI/CD 触发,只按运行时长计费。两者配合形成从探索到生产的自然生命周期。

性能与成本调优的关键是显式声明资源。默认每批给 4 核与 16,000MB 内存,容易造成内存缺或 CPU 浪费。内存密集任务应单独调高 spark.driver.memory 和 spark.executor.memory;计算密集任务则调整 core 数,注意增加核数会按比例附带内存。Google 新推出 history-based autotuning:把重复批任务归为 cohort,基于历史 telemetry 自动给出优化。

还要用 spark.dynamicAllocation.maxExecutors 设上限,作为预算保险丝;SLA 高的任务放高上限,夜间批处理收紧缩,保证 DCU 消耗平稳。shuffle 场景下默认 200 分区,数据大时单分区超内存会溢写到磁盘,建议按总体数据量调整分区数,让每分区承载约 100-200MB。

生产故障排查可直接在 console 调出 Gemini Cloud Assist。官方示例中,批任务报 exit code 1,助手自动分析 driver 日志,指出缺少源码 GCS 路径等运行参数;第二次失败定位到 TypeError,原因是 schema 把金额与 ID 推断成字符串,导致除法崩溃。随后用户让助手生成修复代码,助手返回带 coalesce 和 try_cast 的写法,过滤坏记录而不中断整个批次。

这篇指南按部署选型、调优、AI 排障三部分展开,各部分可独立阅读。适合想减少基础设施负担、同时保留成本与性能控制的数据工程团队。

Serverless Apache Spark on Google Cloud: Architecture & AI Troubleshooting

查看原文