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 流水线等系统触发。任务结束后计算资源释放,成本边界更清晰。
一个实用的开发流程是:
- 使用交互式会话探索数据、确认字段类型并验证转换逻辑。
- 将稳定逻辑整理成参数化的
.py或.jar应用。 - 使用 Serverless batch 执行生产任务。
- 通过编排系统传入输入路径、输出路径和运行日期。
这种方式把“探索时的即时反馈”和“生产时的自动化执行”分开,既方便调试,也避免长期保留开发会话。
用运行时参数控制性能和 DCU 成本
Serverless 并不意味着不需要调优。默认资源规格对简单任务足够,但内存密集型、计算密集型和大规模 Shuffle 作业通常需要显式设置运行时参数。
资源形状要和瓶颈匹配
可以这样判断:
- 如果任务频繁出现 OOM,优先检查
spark.driver.memory和spark.executor.memory。 - 如果 CPU 长时间满载而内存使用率较低,调整
spark.driver.cores和spark.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 分区
groupBy、join 和 distinct 等宽转换会产生 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 自动推断把 amount 或 quantity 当成字符串,而代码却直接进行除法。源文件中还可能混有空字符串、文本金额或零值分母。
排查时应同时确认两件事:转换表达式的输入类型,以及原始文件中是否存在不符合预期的记录。只修改 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,它可以减少闲置资源;对于需要深度生态兼容或底层控制的长期任务,托管集群仍然是更稳妥的选择。