引言EMR Serverless Spark AI Function 自推出以来已在智能智驾、具身智能、互联网等多个行业落地服务了大批生产客户。在此前的内容中我们介绍了 AI Function 的基本用法与核心能力本篇作为延续进一步解读并发控制、算子下推等执行机制并聚焦一个随数据规模增长愈发关键的问题如何控制 AI 任务的整体计算成本。一项大规模 AI 任务的成本主要来自两部分调用模型产生的推理费用以及运行 Spark 作业产生的计算资源费用。围绕这两类成本本文介绍两个作用于不同维度的降本手段——感知 AI Function 的查询优化减少真正进入模型的调用量与 Batch File 异步批量推理以更低的单位价格执行推理并释放等待阶段的计算资源并说明它们各自的适用场景与叠加方式。任务总成本 ≈ 模型推理成本 Spark 计算成本模型推理成本 ≈ 有效调用量 × 单次平均 Token × 模型单位价格Spark 计算成本 ≈ 活跃计算资源规模 × 活跃时长 × 资源单价降本维度代表能力适用范围主要价值减少有效调用量AI 查询优化在线模式与 Batch File 模式过滤无效数据避免重复调用降低模型单价与计算成本Batch File 异步批量推理对时延不敏感的离线任务使用批量价格收缩等待阶段的 Executor大规模 AI 任务的成本为什么容易失控在普通 Spark 查询中过滤、投影、聚合等本地算子的单位成本相对稳定优化器通常关注扫描量、Shuffle、Join 策略和数据倾斜。AI 函数却改变了成本结构一次远程模型调用的延迟和价格可能远高于同一行上的本地表达式计算如果计算节点还需要长时间等待远端结果Spark 资源费用也会随等待时间继续累积。这意味着大规模 AI 任务的成本治理至少要回答三个问题实际需要调用多少次模型、每次调用以什么方式计费以及 Spark 计算资源需要活跃多久。两个逻辑等价的执行计划最终的 AI 调用量可能相差几个数量级同样一批请求在线与 Batch File 不仅面向不同的时延目标和模型价格也会形成不同的计算资源占用方式。例如SELECTreview_id,ai_sentiment(review_text)ASsentimentFROMreviewsWHEREreview_date2026-01-01LIMIT1000;这条 SQL 不会先推理全量再截取 1000 行——Spark 会先用日期条件缩小范围也会在拿到足够结果后停止处理。但在分布式执行中系统仍可能提前准备候选数据导致少量已推理的行未进入最终结果。对本地计算这微不足道对模型调用则意味着多余的 Token 和费用。查询优化会尽量先确定真正需要的 1000 条数据再推理让调用量更贴近最终结果规模。因此适用范围最广的降本动作不是把请求发得更快而是先确认哪些请求根本不需要发送。降本维度一用 AI 查询优化减少调用量AI 查询优化让 Spark 在计划阶段识别 AI 表达式意识到Project 通常是轻量本地计算的假设不再成立从而在保证结果语义的前提下把便宜、确定的数据缩减尽量安排在模型调用之前。下文选取 Filter 前置、Limit 与排序协同、输入去重三类典型场景说明优化思路引擎还内置了更多 AI 查询优化规则能根据实际 SQL 自动协同生效。2.1 复杂表达式下的 Filter 前置先缩小数据再推理大模型调用受采样和服务状态影响同一输入不保证每次结果相同因此 EMR Serverless Spark 按非确定性语义处理 AI Function避免通用规则随意重排模型调用。但这也意味着当普通计算与 AI Function 出现在同一层时原生 Spark 往往不会让 Filter 穿过混合计算即使过滤条件完全不依赖模型结果。下面这个例子中过滤条件依赖新生成的normalized_text但不依赖 AI 返回的sentimentSELECTnormalized_text,sentimentFROM(SELECTupper(review_text)ASnormalized_text,ai_sentiment(review_text)ASsentimentFROMreviews)enrichedWHEREnormalized_textLIKE%REFUND%;AI 查询优化器能识别过滤条件与模型结果之间的依赖关系先生成标准化文本并完成过滤再仅对保留下来的行做情感分析。如果过滤后只剩 20% 的数据AI 调用量也随之缩小到约 20%。反之如果过滤条件引用了 AI 结果如sentiment negative过滤就不能越过对应的 AI 计算。2.2 Limit 与排序协同重新权衡短路与推理成本需要区分单独LIMIT和ORDER BY ... LIMIT两类场景。对于没有排序的LIMIT 1000Spark 本身会在拿到足够结果后停止处理但分布式执行中仍可能提前准备候选数据造成少量额外调用查询优化会尽量先确定真正需要的数据让调用量进一步贴近 1000。带排序的 Limit 优化空间更大。例如按评论时间取最新 1000 条再做情感分析只要排序键不依赖 AI 结果引擎就可以先完成排序与裁剪再只对最终 1000 行调用模型。但如果查询要求按 AI 评分排序后取前 1000 条所有候选行的评分都可能影响排名此时不能把LIMIT提前——查询优化遵循的首要原则仍是保持 SQL 语义。下图以一条带LIMIT 10的查询为例展示了优化前后的执行计划与 AI 指标对比未优化时AIProject位于CollectLimit之前引擎先对 26 行候选数据逐一调用模型再截取最终 10 行开启查询优化后GlobalLimit被推至AIProject之前仅对最终需要的 10 行执行推理。对应地number of AI requests从 26 降至 10number of AI input tokens从 65,572 降至 25,220AI 请求耗时从 1.6 分钟缩短至 28.8 秒。这些指标请求数、Token 用量、限流重试次数等是 EMR Serverless Spark 在原生 Spark UI 基础上新增的 AI 可观测维度用户无需额外埋点即可在作业详情中直接观察优化效果。2.3 输入去重相同问题只问一次真实数据中经常存在大量重复输入例如模板好评、重复工单、批量复制的商品描述。当同一个 AI 函数以相同参数处理相同输入时引擎可以只计算一次再把结果复用到对应行。相同输入指完整的 AI 调用语义一致——函数类型、模型服务、标签集合和推理选项都相同结果才会复用。去重本身有额外的数据组织和结果回填开销因此更适合重复率较高、单次调用较贵的数据集。2.4 不止三个示例而是一套优化体系上述三个场景是优化思路的代表性切面并非完整列表。EMR Serverless Spark 已围绕表达式依赖、数据缩减、结果复用和复杂查询协同等方面构建了一整套 AI 查询优化规则可以单独生效也可以在一条 SQL 中协同工作但都遵循同一条原则在不改变查询语义的前提下让尽可能少的数据到达 AI 计算节点。用户看到的仍是一条普通 SQL不需要为了节省调用而手工拆分任务或维护中间表。查询优化回答了“能不能少调一些”无论后续采用在线还是离线执行都能产生收益。对于必须尽快返回的任务优化后的请求继续走在线模式而对于能够接受非实时完成的任务还可以进一步追问剩下这些调用能否用更具性价比的批量方式完成降本维度二接受非实时用 Batch File 降低模型与计算成本3.1 Batch File 适合解决什么问题阿里云百炼 Batch File API 面向对时延不敏感的大规模推理任务采用异步批处理范式先提交一批请求由服务端在完成窗口内处理再统一获取结果。与实时 API 的发送一个请求、等待一个响应不同它把推理从同步调用变成了后台任务。但把这套异步批处理能力接入 Spark并不只是把多行请求拼成一个文件。批量任务可能需要较长时间才能完成如果 Spark Task 在此期间一直占用 Executor 轮询即使模型调用价格降低等待过程仍会带来持续的计算资源开销。要同时降低两类成本关键是让 Spark 计算资源不必陪着远程任务一起等待。因此对于离线数据管线这种模式会同时影响两张账单模型推理费用更低Batch 推理通常按对应实时模式价格的 50% 计费具体支持范围和价格以百炼最新计费说明为准Spark 计算资源费用更低异步提交完成后Executor 不需要持续占用来等待模型结果在 Serverless 资源策略允许时可以收缩执行计算资源结果就绪后再恢复计算。3.2 异步 Batch 的三个阶段EMR Serverless Spark AI Function 的异步 Batch 执行包含如下三个阶段提交Spark 读取输入数据生成批量请求并提交给模型服务同时保留后续恢复结果所需的关联信息等待提交完成后释放执行计算节点由轻量的任务协调逻辑跟踪批量任务状态回收模型服务完成推理后Spark 恢复计算拉取结果并与原始记录对齐继续执行后续 SQL。这套机制的关键价值不是“让模型更快返回”而是把远端服务的等待时间与 Spark 计算资源解耦。在异步模式下Executor 不需要为了轮询状态持续占用工作空间能否缩容以及最终资源费用仍取决于用户的 Serverless 资源策略和作业中的其他活动阶段。状态恢复、失败处理和结果对齐由引擎统一管理。对用户来说Batch File 仍然只是同一条 SQL 的一种执行方式不需要额外编写上传文件、轮询任务和回填结果的脚本。3.3 在非实时场景中两个降本维度如何叠加查询优化与 Batch File 可以分别使用也可以同时生效。查询优化属于通用能力在线与离线任务都可以受益Batch File 则适用于业务能够接受非实时返回的离线任务。当任务同时具备离线、规模大、成本敏感这三个特征时两项能力可以自然叠加只开启 Batch仍可能把大量本可过滤或去重的数据送进模型先通过查询优化缩减调用再进入 Batch才能同时降低调用量和单位价格。在这种离线场景下执行路径是普通 Spark 算子先完成扫描裁剪、过滤、排序和必要的数据缩减查询优化进一步减少真正需要推理的行和重复输入剩余请求再进入异步 Batch以更低的单位价格处理结果返回 Spark继续参与后续投影、写表和分析。两项能力共同服务于“降本”但适用范围与作用维度不同查询优化负责减少调用量Batch File 在业务接受非实时的前提下进一步降低模型单价和计算资源成本。实战500 万条电商评论的双维降本下面选择一个明确允许非实时完成的离线场景把两项能力串起来。示例数字用于说明计算方法不代表固定的产品性能或价格承诺。4.1 业务背景某电商平台在 Paimon 表中积累了 500 万条用户评论。数据团队希望完成三类标注情感倾向正面、负面、中性问题分类物流、质量、价格、服务等关键实体商品名称、缺陷类型、配送时长等。结果会写入下游分析表用于品类洞察、客户体验分析和客服工单路由。传统 Python 脚本不仅需要逐条调用模型还要自行实现并发控制、限流退避、失败重试、断点续传和结果落盘。数据规模越大这些非业务代码越容易成为主要维护成本。4.2 用户看到的仍然是一条 SQL-- 启用 Batch File并使用可释放 Executor 的异步模式SETspark.emr.serverless.ai.batchFile.enabledtrue;SETspark.emr.serverless.ai.batchFile.modeasync;-- 数据中存在较多模板评论时启用输入去重SETspark.emr.serverless.ai.deduplicate.enabledtrue;WITHenriched_reviewsAS(SELECTreview_id,review_text,to_date(review_time)ASreview_date,lower(trim(product_category))ASnormalized_category,ai_sentiment(review_text)ASsentiment,ai_classify(review_text,ARRAY(物流问题,质量缺陷,价格争议,服务态度,正面好评))ASissue_category,ai_extract(review_text,ARRAY(product_name,defect_type,delivery_days))ASentitiesFROMods_user_reviews)INSERTINTOreview_labelsSELECTreview_id,normalized_categoryASproduct_category,review_text,sentiment,issue_category,entitiesFROMenriched_reviewsWHEREreview_date2026-01-01ANDnormalized_categoryIN(electronics,home_appliance);这条 SQL 在 CTE 中完成日期转换、品类标准化和三项 AI 标注后续 Filter 只依赖普通派生列执行模式、恢复流程和并发细节由引擎统一管理。4.3 调用量是怎样降下来的假设示例数据符合以下分布阶段数据量说明原始评论500 万行混合计算的全量候选数据日期与品类过滤后120 万行查询优化将不依赖模型结果的 Filter 前置去除重复 AI 输入后85 万条唯一输入模板评论只推理一次三个 AI 函数的最终调用量255 万次85 万 × 3在这个混合表达式场景中如果没有查询优化包含非确定性 AI 调用的整体计算无法被简单拆分500 万行数据都可能先进入三个 AI 函数潜在调用量为 1500 万次。查询优化在确认过滤条件不依赖模型结果后先将数据收缩到 120 万行再结合输入去重收敛到 85 万条唯一输入。三个 AI 函数最终只需约 255 万次调用降到未优化基线的 17%即减少约 83%。4.4 两个降本维度如何叠加为了避免不同模型之间的价格差异影响比较下面采用相对成本进行估算查询优化后的调用量系数 255 万 / 1500 万 17%Batch 单位价格系数 50%相对推理成本 17% × 50% 8.5%在单次平均 Token 基本相同的前提下组合方案的模型推理费用约为“全量实时逐行调用”的 8.5%对应约 91.5% 的理论降幅。上述估算仅覆盖模型推理费用Spark 计算侧的收益取决于资源规格与伸缩策略需结合实际作业单独评估。4.5 方案对比维度自建 Python 调用脚本在线 AI Function QO异步 Batch File QO开发方式自行编排 API 与状态SQL / DataFrame APISQL / DataFrame API 执行配置示例 AI 调用量取决于脚本是否主动过滤与去重约 255 万次约 255 万次模型单位价格实时价格实时价格通常为对应实时价格的 50%远端等待期脚本进程持续管理Executor 参与在线执行异步模式可释放 Executor限流与重试用户自行实现引擎统一处理引擎统一处理任务恢复与结果对齐用户自行实现引擎统一处理引擎统一处理适合场景小规模定制流程低延迟、交互式或持续到达的数据大规模、离线、成本敏感任务如何开启与调优白名单说明Batch File 异步批量推理功能当前处于白名单测试阶段如需使用请联系 Serverless Spark 团队申请开通。5.1 最小配置如果目标是使用异步 Batch 并在等待阶段释放 Executor需要同时开启 Batch File 和异步模式SETspark.emr.serverless.ai.batchFile.enabledtrue;SETspark.emr.serverless.ai.batchFile.modeasync;如果数据重复率较高再开启输入去重SETspark.emr.serverless.ai.deduplicate.enabledtrue;5.2 按 AI Function 选择执行模式会话级配置适合整段作业采用统一策略。如果只希望为某个 AI Function 开启 Batch File或者需要为不同处理阶段分别选择 Batch File 的同步、异步模式也可以通过函数的options参数设置batch_mode-- 为该 AI Function 选择异步 Batch File 模式SELECTreview_id,ai_sentiment(review_text,options{batch_mode:async})ASsentimentFROMreviews;batch_mode可以设置为sync或async设置为true时则启用 Batch File 并沿用会话默认模式。实际使用时建议同一个SELECT投影中的多个 AI Function 保持一致的batch_mode如果它们的时延要求不同可以拆分为不同处理阶段。多模态数据当前不支持 Batch File当前 Batch File 模式暂不支持多模态数据包括携带图片等多模态输入的ai_query和ai_embedding_multimodal。这类调用会使用在线模式如果它们与文本 AI Function 出现在同一个SELECT投影中该投影也会整体按在线模式执行。建议将文本离线推理与多模态处理拆分为不同阶段。5.3 常用参数-- 单个批量文件的最大请求数当前默认值为 10000SETspark.emr.serverless.ai.batchFile.maxRequestsPerFile10000;-- 批量任务完成窗口当前默认 24h可在支持范围内调整SETspark.emr.serverless.ai.batchFile.completionWindow24h;参数并不是越大越好。文件过大会拉长单个批次的恢复粒度过小则增加任务和文件数量完成窗口应根据业务 SLA 选择。5.4 PySpark DataFrame APIfromemr_serverless_spark_aiimportai_sentiment,ai_classifyfrompyspark.sql.functionsimportcol spark.conf.set(spark.emr.serverless.ai.batchFile.enabled,true)spark.conf.set(spark.emr.serverless.ai.batchFile.mode,async)spark.conf.set(spark.emr.serverless.ai.deduplicate.enabled,true)reviewsspark.table(ods_user_reviews)resultreviews.select(col(review_id),ai_sentiment(col(review_text)).alias(sentiment),ai_classify(col(review_text),[物流,质量,价格,服务],).alias(category),)result.write.mode(overwrite).saveAsTable(review_labels)SQL 与 DataFrame API 共用同一套执行与优化能力团队可以按现有工程习惯选择接口。选型建议什么时候用在线什么时候用 BatchBatch File 不是在线模式的替代品两者面向不同的延迟目标。判断问题更适合在线模式更适合异步 Batch File用户是否正在等待结果是否是否要求秒级或分钟级响应是否允许完成窗口数据是否持续到达流式或持续到达有明确批次边界任务规模小到中等、重视即时性大规模、重视吞吐与成本计算资源是否可在等待期缩容通常不关注希望释放 Executor典型场景交互式分析、在线服务、实时辅助T1 打标、历史回填、离线 ETL、周期性内容处理还有两个容易忽略的边界查询优化不只服务于 Batch。Filter、Limit 等调用削减能力对在线模式同样有效并非所有输入形态都适合 Batch。例如多模态请求涉及媒体读取和有效期管理当前会使用在线执行。具体函数与模型支持范围应以实际测试为准。总结EMR Serverless Spark AI Function 把大模型能力带入用户熟悉的 SQL 与 DataFrame 工作流同时让 Spark 能够理解 AI 调用是一种昂贵、需要治理的外部计算。面对不断增长的数据规模降本需要同时关注模型服务和 Spark 计算两张账单调用多少次、以什么价格调用以及计算资源需要活跃多久。本文介绍的两个降本维度分别解决了不同问题也有不同的适用边界AI 查询优化通过 Filter、Limit、输入去重等手段在计划阶段减少不必要的模型调用在线与离线任务都可以受益Batch File 异步批量推理面向能够接受非实时完成的离线任务用更具性价比的执行方式处理请求并将远端等待与 Executor 资源占用解耦。对于时延敏感任务在线 AI Function 加查询优化就是完整方案对于成本敏感的离线任务则可以在此基础上选择 Batch File形成“先减少调用再降低模型单价与计算资源成本”的组合优化路径。对数据团队而言真正有价值的不只是把几段 API 调用改写成一条 SQL而是把过去散落在脚本中的并发、限流、恢复、成本和资源问题收敛为平台能够统一优化和治理的数据处理能力。这也是 Serverless Spark 从“大数据计算引擎”走向“AI 原生数据处理平台”的关键一步。