返回
RSS Google Cloud Data Analytics Blog AI 逐段翻译 发布 2026-08-19 08:00 收录于 08-28

Google Cloud Serverless Spark架构与AI故障排查

DataHot 速览

Google Cloud 官方博客介绍了 Serverless Apache Spark 的两种执行模式:交互式会话适合迭代探索,批处理用于自动化生产,并说明了各自的计费特点。文章还展示了从开发到生产的生命周期切换方式,以及高级性能调优和 DCU 成本优化内容。

为什么值得关注:数据平台从业者可借鉴 Serverless Spark 在不同阶段的选型与成本控制方法,并了解 Google Cloud 的 AI 辅助排障实践。

本文目录 9 节
  1. 无服务器交互式会话
  2. 无服务器批处理
  3. 开发到生产的生命周期
  4. 第 2 部分:高级性能调优和 DCU 成本优化
  5. 自定义驱动器和执行器形状
  6. 控制自动缩放边界
  7. 管理 shuffle 存储效率
  8. 第 3 部分:使用 Gemini Cloud Assist 进行运维诊断
  9. 阶段 1:诊断缺失的执行参数

译文

AI 逐段翻译

一旦选择无服务器部署模式,您必须根据开发阶段和运营要求选择合适的执行模型。Apache Spark 托管服务提供两种运行无服务器工作负载的选项:

无服务器交互式会话

交互式会话非常适合迭代和探索性用例。您可以编写代码块、检查中间 DataFrame、修改变量,并在数据集热保存在内存中的情况下生成可视化。

  • 主要接口:专为人机交互设计。开发人员使用他们选择的 IDE(如 Colab、Gemini Enterprise Agent Platform Workbench、Antigravity、Jupyter 笔记本等)逐单元格执行代码。
  • 空闲成本概况:计算资源保持活跃以支持开发人员思考期间的即时执行,如果会话处于非活动状态,可能会产生一些空闲计算费用。

无服务器批处理

当您知道要运行什么,并且需要自动化、非交互式执行时,批处理非常有用。引擎从头到尾运行完整、打包的 PySpark 脚本(.py)或 Java/Scala 应用程序文件(.jar),无需人工干预。

开发到生产的生命周期

这些执行选项旨在协同工作,形成自然的流水线生命周期。在初始开发阶段,您在笔记本界面中打开无服务器交互式会话,以探索数据集、清理模式并原型化转换。一旦逻辑得到验证且转换最终确定,您将代码打包为 Python 脚本,并将其调度为由 Apache Airflow 托管服务编排的无服务器批处理作业,用于生产执行。这种过渡在保持运营可靠性的同时,最大限度地降低了持续开发成本。

第 2 部分:高级性能调优和 DCU 成本优化

虽然无服务器托管 Spark 消除了集群维护的运营开销,但在生产企业级流水线上使用默认设置可能会导致性能瓶颈或预算浪费。必须使用运行时配置属性在提交时明确声明资源分配,以维持高效的数据计算单元(DCU)消耗率。

Google 最近引入了基于历史记录的自动调优。在无服务器的上下文中,此功能自动基于最佳实践和历史执行应用优化。它通过将重复的批处理工作负载分组到 Google 所称的 cohort 中来实现。自动调优器分析同一 cohort 名称下先前运行的遥测和统计信息,以找出瓶颈所在。

自定义驱动器和执行器形状

默认情况下,无服务器批处理分配通用规格(4 核和 16,000MB RAM)。这可能会导致关键效率问题,具体取决于应用程序的性质:

  • 内存密集型作业:处理高度未压缩数据量的流水线可能会遇到内存不足(OOM)错误并崩溃。为此,请使用 spark.driver.memory 和 spark.executor.memory 独立增加堆大小。
  • 计算密集型作业:运行数学建模或繁重分词等处理密集型作业可能会使 CPU 饱和,而昂贵的 RAM 则闲置。通过明确调整 spark.driver.cores 和 spark.executor.cores 来微调每个实例的处理并发性。

请记住,默认情况下,增加核心数会自动按比例配置内存,以匹配 vCPU 与 RAM 的比率。这就是为什么覆盖核心数和内存的值至关重要。

控制自动缩放边界

托管 Spark 无服务器根据积压任务动态扩展活跃执行器的数量。然而,如果引入恶意代码循环或未优化的笛卡尔联接,不受约束的扩展可能导致预算超支。

作为防御性护栏,始终使用 spark.dynamicAllocation.maxExecutors 声明明确的上限。这充当您的预算安全开关。通过将其上限设置为合理的上限,您保证即使代码行为不佳,作业也不会超过固定的基础设施足迹。

  • 高优先级(SLA 驱动):将 maxExecutors 设置为较高的上限,以允许资源突发并最小化整体运行时间。
  • 低优先级(夜间批处理):将 maxExecutors 设置为较低、紧凑的上限。工作负载将运行更长时间,但将消耗可预测、平稳、成本高效的 DCU 流。

管理 shuffle 存储效率

当执行涉及宽转换(如 groupBy()、join() 或 distinct())时,数据必须在网络间重新分配,生成称为 shuffle 存储的中间磁盘写入

Spark 默认静态设置为 200 个分区(spark.sql.shuffle.partitions)。如果您正在处理巨大的多 GB 数据集,200 个分区意味着每个单独块将太大。当分区大小超过可用执行器 RAM 时(例如,1GB 分区试图在分配的 0.5GB 堆空间内处理),数据会溢出到磁盘。这会减慢执行速度,并为高级或标准 shuffle 存储块产生额外计费。一个有用的经验法则:根据总数据大小动态调整分区参数,使每个分区在内存中处理大约 100MB 到 200MB 的数据。这可能需要几次迭代才能达到最佳结果。

上述属性是主要的可调属性。其他无服务器运行时配置属性可在此链接

第 3 部分:使用 Gemini Cloud Assist 进行运维诊断

当自动化数据管道在生产环境中失败时,数据工程师传统上被迫花费数小时筛选驱动程序和执行器中冗长、不连贯的日志文件。Apache Spark 托管服务通过原生集成Gemini Cloud Assist到 Google Cloud 控制台中,使工程师能够使用自然语言诊断和解决故障,从而解决了这一痛点。

为了说明这一运维转变,我们考察一个失败的 PySpark ETL 管道的典型故障排除生命周期,该管道从Google Cloud Storage(GCS)存储桶读取客户交易数据,应用转换,并遇到意外的运行时错误。

阶段 1:诊断缺失的执行参数

在新管道的首次执行尝试期间,批处理作业状态从“待处理”变为“运行中”,最终以通用退出消息“应用程序以退出代码 1 失败”结束为失败状态。

工程师无需手动查询 Cloud Logging 或在控制台的多个部分之间导航,即可找到错误日志并选择“调查日志”选项。此操作会打开一个原生对话窗格,Gemini Cloud Assist 会在其中自动分析驱动程序遥测和系统日志。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存