Spark调用大模型实战:mapPartitions并发控制与工程化落地

发布时间:2026/9/10 10:42:35
Spark调用大模型实战:mapPartitions并发控制与工程化落地 如果你在数据团队里待过一阵子大概率会遇到这种需求给几百万条文本打标签、抽取实体、做情感分析或者把用户聊天记录批量生成模型训练样本。业务方开口就是“用大模型跑一下就行”但真到动手阶段才发现把大模型接进 Spark 里并没有想象中那么无脑。最常见也最容易被吐槽的写法是在 UDF 里直接requests.post调大模型 API。跑小数据量看不出问题一旦上了千万级数据连接打满、超时重试混乱、Executor 被拖垮各种问题一起冒出来。这篇文章我想从工程视角把“Spark 调用大模型”这件事完整拆开先说清楚哪些场景适合做、哪些场景其实不适合然后对比几种主流的接入方案给出选型建议再放出可以直接改着用的 PySpark 代码把并发、超时、重试、序列化这些硬骨头一个一个讲透。适合读这篇文章的人我默认是两类一类是做数据平台、离线数仓、算法工程的开发正在琢磨怎么把大模型能力批量落到数据管道里另一类是准备数据开发或大数据架构面试的同学因为这个问题这两年出现频率越来越高面试官问的不是你会不会调 API而是你能不能把分布式计算和大模型服务之间的资源、并发、容错关系讲明白。我会尽量把代码、参数、踩坑记录都放出来文章里的路径都是我实际跑过的照抄可以少走很多弯路。1. 为什么会在 Spark 里调用大模型场景与挑战1.1 典型应用场景先说清楚一件事用 Spark 调大模型本质上是“批量推理”不是“在线推理”。在线推理要求毫秒级响应那是后端服务的事跟 Spark 没什么关系。Spark 的定位是离线批处理所以凡是“有一批存量数据想用大模型处理一遍”的场景都适合拿到这里来讨论。我实际见到过并且自己参与过的场景大概有这么几类文本批量打标与分类比如把平台上积累的历史工单、评论、投诉内容统一做一次意图分类量级通常在百万到千万条。用规则或小模型效果不好用大模型跑一遍就能把标签体系更新一遍。信息抽取与结构化从非结构化文本里抽出时间、地点、金额、人物等实体或者从长文档里抽取摘要。这个场景大模型比传统 NER 模型灵活得多尤其适合字段不固定、规则经常变的抽取任务。离线生成训练样本用大模型给原始数据生成蒸馏样本、构造正负例、做数据增强。比如先用大模型生成一批高质量的问答对再拿去微调小模型这是现在很多团队降本增效的标准路径。内容审核与风险识别把待审核的数据批量过一遍大模型输出风险等级、违规类别。这类任务往往要和规则引擎结合Spark 负责把规则命中的部分再交给大模型做二次判断。知识库向量化前置处理虽然向量化通常用小模型做但在构建 RAG 知识库时很多团队会先用大模型对文档做清洗、切分、标题生成、摘要再做 embedding。这些场景有一个共同点输入侧数据规模大输出侧对实时性要求不高但要求在合理时间内跑完并且要能监控进度、控制成本、失败能重试。1.2 核心矛盾与常见误区Spark 调用大模型本质上是在两个完全不同的系统之间做桥接。Spark 是分布式批处理框架数据分布在很多 Executor 上并行算大模型是外部服务不管是公有云 API 还是私有化部署的推理服务它有自己的并发上限和时延。两者之间最核心的矛盾有三个第一个是连接与并发模型不匹配。Spark 的并行度可以开到成百上千一个 Executor 上可以同时跑很多 Task。如果你在每个 Task 里都新建一个 HTTP 连接去调 API很快会把服务端连接池打满也会触发限流。更麻烦的是Python UDF 是逐行调用的每一行都做一次完整的 HTTP 往返连接无法复用性能会非常难看。第二个是资源模型的冲突。大模型推理需要 GPU而 Spark 的 Executor 通常跑在 CPU 容器里。就算你用的是 Spark 3.4 之后支持 GPU 调度的版本真正跑大模型推理时GPU 显存和 Executor 堆内内存的管理也是两套逻辑。多数情况下让 Spark 做数据准备和结果回填把推理压力放在独立服务上才是合理的分工。第三个是容错哲学完全不同。Spark 对 Task 失败的态度是“重新算一遍”某个分区失败了就重新执行该分区。但大模型调用是有外部副作用的Token 已经消耗了API 可能已经返回了结果只是网络超时了。如果你很天真地开启 Spark 重试机制同一个 prompt 可能被重复发送很多次账单会很难看。所以在 Spark 里调用大模型重试逻辑必须自己控制不能依赖 Spark 的 Task 重试。很多人问我第一个问题就是“直接写个 UDF 调一下不行吗”行数据量小确实行我就这么干过。但“能跑”和“能用”是完全两码事。如果不控制并发、不复用连接、不做超时管理百万级数据基本必挂。下面我会把不同方案摆出来对比你就能明白为什么工程上最终都走向了“客户端池化 分区处理”这条路。2. 方案选型对比四种主流调用思路2.1 方案AUDF 直连 HTTP API这是最容易上手、也最容易翻车的方案。直接在 PySpark 里注册一个 UDFUDF 内部用requests或 OpenAI 兼容 SDK 调大模型 API然后df.withColumn(result, call_llm_udf(col(text)))一把梭。from pyspark.sql.functions import udf from pyspark.sql.types import StringType import requests def call_llm(text: str) - str: resp requests.post( https://api.example.com/chat/completions, json{model: demo-model, messages: [{role: user, content: text}]}, timeout30 ) data resp.json() return data[choices][0][message][content] call_llm_udf udf(call_llm, StringType())这个方案在小数据量验证时非常方便几千条数据跑一遍完全没问题。但问题也很直接每行数据都新建一次 HTTP 连接TCP 握手开销巨大而且很多 SDK 内部会自己管理连接池放在 UDF 里每行新建客户端连接池根本来不及复用。Spark 默认的并行度如果比较高会在同一瞬间发起大量并发请求直接把服务端打挂或者触发 429 限流。超时和重试逻辑写起来很别扭。timeout一设模型生成时间长一点就超时不设超时某个请求卡住整个 Spark 任务可能挂在最后几个 Task 上特别难受。我的建议这个方案只用来做连通性验证或百条级以下的试验不要用于生产。如果你判断数据量短期内都不会超过几千条那用这个方案也无可厚非但最好还是提前预留改造空间。2.2 方案BmapPartitions 客户端池化这是目前工程上最主流、性价比最高的方案。核心思路很简单不用 UDF 逐行调用而是在每个分区内部初始化一次客户端然后在这个分区内用并发线程池去调用大模型。为什么用mapPartitions因为 Spark 中数据是以分区为单位的一个分区内的数据会在同一个 Executor 上顺序或并行处理。在分区级别做初始化客户端只需要创建一次连接可以复用同时你可以在分区内部用线程池控制并发把并发量稳定在一个服务端能接受的范围内。这段逻辑用代码表达大致是这样的框架def process_partition(iterator): # 每个分区只会执行一次初始化 client LLMClient(...) for row in iterator: yield do_call(client, row)这里有几个容易踩的细节如果你在 UDF 或者map里创建客户端客户端会被序列化复制到每个 Task 里要么报Task not serializable要么每个 Task 都重新建连接性能还不如方案A。而在mapPartitions里创建客户端只属于当前分区所在的 Python 进程一切都在本地不涉及序列化问题。我在后面的实战章节会给出完整的可运行实现这里先不展开。2.3 方案C外部服务化 消息队列 结果回填当数据量到了日均千万甚至亿级或者大模型服务不稳定、需要独立扩缩容的时候方案B也会遇到瓶颈。这时更稳妥的做法是把 Spark 作业拆成两个阶段Spark 负责数据预处理和下发任务独立推理服务负责消费和执行结果再回填到表里。具体的流程一般是这样的Spark 读取原始数据清洗、裁剪、拼 prompt把任务写入消息队列Kafka 或 RocketMQ消息体里带上业务主键 ID 和 prompt。独立的推理服务集群从消息队列消费任务调用大模型 API 或本地推理服务得到结果后写回结果表。Spark 在第二阶段或者另一个周期调度扫描结果表把有结果的记录和原始表关联完成后续加工。这个方案把 Spark 和大模型服务之间的耦合彻底解开了。推理服务可以按照模型吞吐独立扩缩容队列可以在服务端抖动时起到削峰填谷的作用Spark 不需要关心调用大模型的细节只关心数据怎么写、怎么读。调用的吞吐上限从“Spark 分区并发数”变成了“消息队列的消费能力”扩展性好了很多。代价也很明显架构复杂度增加了。你需要维护消息队列、消费服务、结果表的状态管理还要处理消息丢失、重复消费、结果回填延迟等问题。我的判断是如果你的数据量单次任务在百万到千万条量级用方案B就够了超过这个量级且要常态化跑再考虑方案C。不要一上来就上消息队列很多团队就是被自己的过度设计拖垮的。2.4 方案D本地 GPU 推理服务 Spark 协同这一两年本地模型部署工具成熟得很快VLLM、TGI、Ollama 这些项目已经能把开源模型部署成 OpenAI 兼容的 HTTP 服务。于是有一种新的玩法Spark 集群边上放一台或几台 GPU 服务器用 VLLM 部署开源模型Spark 在内网直接调用这个服务。这个方案最大的优势是数据不出域。对于金融、医疗、政务这类对数据出境有严格要求的场景公有云 API 往往不能直接用本地部署几乎是唯一选择。而且如果推理量足够大本地部署的单 Token 成本可以远低于按量付费的公有云 API长期跑很划算。如果你用的是像 NVIDIA DGX Spark 这样的小型 GPU 工作站在上面跑 Spark 做数据预处理、再驱动部署在本地或同网段的模型服务整个链路都能留在内网效率和安全性都不错。但要注意本地部署大模型不等于没有成本。GPU 服务器的采购或租赁费用、运维成本、模型效果和公有云大模型的差距都要提前评估。我的经验是先用公有云 API 把效果验证清楚再决定要不要本地化。如果模型效果不行省再多的推理成本也没意义。从 Spark 代码的角度看方案D和方案C中的“调用大模型”部分没有本质区别都是 HTTP 调用只是 API 地址从公网变成了内网鉴权从 API Key 变成了内网白名单超时和限流策略需要根据本地服务的实际情况重调。2.5 方案对比与选型建议把这四个方案放在一起看更清楚方案实现复杂度吞吐能力可靠性适用规模适用场景UDF 直连低低差千条级以下实验验证、POCmapPartitions池化中中高中百万到千万级离线批量推理的主流选择MQ服务化高高高千万级以上、常态化核心链路、需要独立扩缩容本地模型Spark中高中高中高视 GPU 规模而定敏感数据、长期大规模推理选型的时候我的思路很简单先看数据量和频率再看数据能不能出域最后看团队运维能力。数据量不过百万、跑一次拉倒的方案B足矣数据量上千万且每天都要跑的直接考虑方案C数据不能出域或者成本敏感的赶紧测方案D。方案A只适合用来验证“这条路能不能走通”。3. 实战PySpark 封装大模型调用的核心实现3.1 环境准备与依赖先交代一下实践环境。我用的是 Spark 3.4 以上版本PySpark 作为开发语言。之所以选 Python 而不是 Scala原因很现实现在大模型相关的 SDK 基本都是 Python 生态requests、openai、httpx这些库在 Python 里最好用而且做数据清洗、prompt 拼接这类的逻辑Python 写起来比 Scala 快得多。你需要提前确认环境里有这些依赖pip install pyspark requests openai其中openai库不是必须的用requests也能实现但如果你调的是 OpenAI 兼容接口用官方 SDK 处理重试、流式这些会省不少事。我在实际项目中两种都写过后面统一用requests做示例因为它能把底层逻辑暴露得更清楚方便排查问题。还需要准备一份环境变量或配置文件保存 API Key、接口地址、模型名。注意 Spark 跑在 YARN 集群上时Python 进程的启动环境和提交任务的客户端环境不一定一样API Key 最好不要直接写死在代码里。我习惯的做法是在spark-submit的时候用--conf spark.yarn.appMasterEnv.LLM_API_KEYxxx传进去代码里用os.environ.get(LLM_API_KEY)读取。spark-submit \ --master yarn \ --deploy-mode cluster \ --conf spark.yarn.appMasterEnv.LLM_API_KEYsk-xxx \ --conf spark.yarn.appMasterEnv.LLM_BASE_URLhttps://api.example.com/chat/completions \ --py-files llm_client.py \ job.py3.2 关键代码mapPartitions 并发调度下面这段代码是我在项目里精简后的通用框架核心思路是每个分区初始化一个客户端分区内部用ThreadPoolExecutor控制并发每条记录里带上业务主键最后返回带主键的结果。from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType import requests import os import time import json import threading from concurrent.futures import ThreadPoolExecutor, as_completed class LLMClient: def __init__(self, api_key, base_url, model): self.session requests.Session() self.session.headers.update({ Authorization: fBearer {api_key}, Content-Type: application/json, }) self.base_url base_url self.model model def chat(self, prompt, max_tokens512, temperature0.0): payload { model: self.model, messages: [{role: user, content: prompt}], max_tokens: max_tokens, temperature: temperature, } resp self.session.post(self.base_url, jsonpayload, timeout(10, 120)) resp.raise_for_status() data resp.json() return ( data[choices][0][message][content], data.get(usage, {}).get(prompt_tokens, 0), data.get(usage, {}).get(completion_tokens, 0), ) def process_partition(iterator, api_key, base_url, model, max_workers8): client LLMClient(api_key, base_url, model) idle_thread threading.local() def do_call(row): text_id, text row start time.time() result, prompt_tokens, completion_tokens client.chat(text) latency_ms int((time.time() - start) * 1000) return (text_id, result, prompt_tokens, completion_tokens, latency_ms) with ThreadPoolExecutor(max_workersmax_workers) as executor: futures [executor.submit(do_call, row) for row in iterator] for future in as_completed(futures): yield future.result()这段代码里我刻意加了一个idle_thread变量但它其实没用到。真正要注意的是ThreadPoolExecutor的并发数和分区的数据量要匹配。如果单个分区的数据量非常大把几万条任务一次性submit进线程池虽然as_completed会边完成边消费但 futures 列表本身会占内存极端情况下也会把 Python 进程内存打爆。更稳的写法是分批提交比如每批 200 条def process_partition(iterator, api_key, base_url, model, max_workers8, batch_size200): client LLMClient(api_key, base_url, model) batch [] with ThreadPoolExecutor(max_workersmax_workers) as executor: for row in iterator: batch.append(row) if len(batch) batch_size: futures [executor.submit(do_call, r) for r in batch] for f in as_completed(futures): yield f.result() batch.clear() if batch: futures [executor.submit(do_call, r) for r in batch] for f in as_completed(futures): yield f.result()这样既能控制并发又能防止 futures 列表无限膨胀。3.3 超时、重试与限流策略大模型 API 的响应时间波动非常大。同样的模型prompt 短的时候一两秒返回prompt 长或者服务端排队的时候可能几十秒。所以超时设置一定要分开连接超时短一点10 秒就够读取超时长一点我一般设 120 秒因为大模型生成一个长回答确实需要时间。requests的timeout参数可以传元组分别指定连接超时和读取超时resp self.session.post(self.base_url, jsonpayload, timeout(10, 120))重试逻辑千万不要交给 Spark 去做而是要在客户端层面自己实现。我建议的重试策略是429限流必须重点试但要退避。500、502、503、504服务端异常重点试同样退避。400、401、403、422客户端错误不需要重试重试多少次都会失败直接抛出异常或者记录错误行。网络异常连接断开、DNS 解析失败可以重试但要设置最大次数上限。一个通用的指数退避重试写法import time import requests def call_with_retry(client, prompt, max_retries3): for attempt in range(max_retries): try: return client.chat(prompt) except requests.exceptions.HTTPError as e: status_code e.response.status_code if status_code in (429, 500, 502, 503, 504): sleep_time min(2 ** attempt, 30) time.sleep(sleep_time) continue raise except requests.exceptions.ConnectionError: time.sleep(min(2 ** attempt, 30)) continue raise RuntimeError(f调用大模型失败已重试 {max_retries} 次prompt 前缀: {prompt[:50]})注意2 ** attempt是 1、2、4、8 的指数退避封顶 30 秒。重试次数一般 3 次就够因为超过 3 次还失败说明服务端可能已经挂了很长时间硬重试只会让问题更严重。我见过有人把重试次数设成 10 次结果服务端故障恢复后所有积压的任务同时重试直接把服务端打挂了。这就是典型的“好心办坏事”。关于限流最核心的一点是并发数要和服务端对齐而不是自己拍脑袋。比如模型服务支持每分钟 600 个请求那你的 Spark 任务总并发最好控制在 10 左右10 QPS因为大模型单次请求的耗时通常在 2-5 秒10 并发一分钟大约能完成 200-300 个请求还有富余。如果开 50 并发看起来吞吐上去了实际上一秒内打出去 50 个请求服务端直接 429然后你的重试逻辑开始退避真实吞吐反而更低。3.4 Schema、序列化与结果对齐调用完大模型后数据的产出不仅仅是结果文本最好把 Token 用量、耗时也一起带出来后续做成本分析和性能调优都离不开这些数据。我推荐的结果 schema 至少包含这几个字段schema StructType([ StructField(id, StringType(), True), StructField(result, StringType(), True), StructField(prompt_tokens, LongType(), True), StructField(completion_tokens, LongType(), True), StructField(latency_ms, LongType(), True), ])然后在 Spark 里这样应用from pyspark.sql import SparkSession spark SparkSession.builder.appName(spark_llm_inference).getOrCreate() df spark.createDataFrame([(1, 这是一条测试文本), (2, 再试一条)], [id, text]) result_rdd df.rdd.mapPartitions( lambda iter: process_partition(iter, api_key, base_url, model) ) result_df spark.createDataFrame(result_rdd, schemaschema)这里有一个很隐蔽的坑df.rdd.mapPartitions返回的 RDD 顺序是不保证的尤其是多分区并行执行时结果顺序和原始顺序完全不同。所以一定要在数据里保留业务主键最后通过join把结果关联回原表。如果你用zipWithIndex给每一行加自增 ID注意要先把 DataFrame 转成 RDD 再操作或者用monotonically_increasing_id()生成唯一递增 ID但后者只保证在同一个 DataFrame 内唯一且递增跨执行阶段不保证连续。如果你对输出格式要求更严格也可以注册成 Spark SQL 的 UDF 直接用但要接受并发的瓶颈。我的建议是能用mapPartitions就不要用 UDF。UDF 是逐行调用的每一行都要经过 Python 和 JVM 之间的序列化边界性能损耗非常大。而mapPartitions是在分区级别处理数据在 Python 进程内批量流转序列化开销少了一个数量级。4. 常见问题与排查实录4.1 Spark on YARN 为什么 CPU 只能用 1 个这几乎是所有 Spark 新手在跑大模型任务时都会遇到的问题提交到 YARN 上的作业日志里显示每个 Executor 只有 1 个 vCore任务并行度上不去调用大模型的吞吐自然惨不忍睹。排查思路要分清是提交参数问题还是集群配置问题。最常见的原因是spark.executor.cores没有设置Spark 默认值是 1。在 YARN 模式下如果你不显式指定每个 Executor 容器就只申请 1 个 vCore。解决办法是在提交命令里加上--conf spark.executor.cores4 \ --conf spark.executor.memory8g \ --conf spark.sql.shuffle.partitions200另一个容易被忽略的原因是 yarn 调度器的限制。yarn.scheduler.maximum-allocation-vcores如果设置成了 1你在 Spark 里申请再多核也没用YARN 会直接按 1 核分配容器。这时候要改的是 YARN 的capacity-scheduler.xml或fair-scheduler.xml配置而不是 Spark 参数。此外如果开了 Dynamic Allocationspark.dynamicAllocation.maxExecutors需要足够大否则 Executor 数量也上不去。到了 Python 场景还有一个 Python 特有的混淆点spark.task.cpus默认是 1表示每个 Task 占用的 CPU 核数。如果你把这个参数设成了 2 或更高每个 Task 会占用多个核Executor 的并行任务数反而下降看起来就是“明明有 4 个核但只有一个任务在跑”。调用大模型这种 IO 密集型任务spark.task.cpus保持 1 就好真正的并发控制在 Python 线程池里做这样 CPU 和网络吞吐都能打满。4.2 Executor 任务不可序列化问题用方案 B 写代码的时候最容易遇到的一个异常是org.apache.spark.SparkException: Task not serializable。原因一般有两个第一个原因是在闭包里引用了无法序列化的对象。比如你在 Driver 上创建了一个 HTTP 客户端然后直接在mapPartitions或者map里用它Spark 会尝试把这个客户端序列化后发给 Executor。requests.Session里包含线程锁和连接池根本没法序列化所以直接报错。解决办法就是不要在闭包外创建客户端把初始化逻辑放到每个分区内部也就是process_partition的开头。第二个原因是闭包里引用了整个 SparkSession 或者一个很大的 DataFrame。有些同学会把 DataFrame 变量写在闭包里Spark 尝试把整个 DataFrame 的逻辑计划序列化下发轻则性能极差重则直接报序列化错误。解决办法是把需要用的值提取成基础类型字符串、整数、字典再传给闭包函数。排查这个问题有个技巧看异常的堆栈信息它会明确指出是哪个对象无法序列化。如果堆栈里提到了LLMClient那基本就是客户端初始化位置不对如果提到了SparkSession或者DataFrame那闭包里多半不小心引用了不该引用的对象。4.3 429 限流与 5xx 服务端错误调大模型 API 时429 Too Many Requests是高频错误尤其是跑批任务时。这个状态码表示你请求太频繁服务端已经限制你了。很多人的第一反应是加大重试次数其实这完全搞反了。429 的正确处理方式是降低并发不是增加重试。我踩过一次很经典的坑。某个任务配了 16 并发跑了几分钟就大量 429重试逻辑开始指数退避退避之后又叠加新请求服务端的限流计数器一直被顶在最大值结果整个任务跑了三个小时还没跑完Token 消耗倒是不少。后来把并发降到 4同样数据量四十分钟跑完了。这个故事的核心教训是大模型服务的并发能力是有限的你只能让 Spark 去适配它而不是让服务端来适配你的 Spark 并行度。遇到 5xx 错误时情况不太一样。500、502、503、504 通常意味着服务端正在经历故障或过载。这时候可以重试但要给服务端留出恢复时间。我建议 503 和 504 的退避时间至少 5 秒起步最多 60 秒如果连续 3 次重试仍然失败就把这批数据标记为失败写入错误表而不是无限重试卡住整个作业。4.4 Executor OOM 与内存模型调优Spark 调用大模型时的 OOM有一大半跟模型输出结果太大有关。大模型单次返回可能包含几千甚至上万 Token 的文本如果每条记录的 prompt 本身就长结果也长Python 进程里有成百上千条这样的数据排队等待处理内存很容易被打爆。要理解这个问题的本质得稍微提一下 Spark 的内存模型。JVM 堆内存被划分为 Reserved、User Memory、Spark Memory又分 Execution Memory 和 Storage Memory等区域。Spark SQL 的 DataFrame 操作主要使用 Execution Memory而 Python UDF 的进程内存并不完全受 Spark 的堆内存管控它由spark.executor.pyspark.memory参数单独限制。如果你跑的是 PySpark 调用大模型可能 OOM 的其实是 Python 进程而不是 JVM。配置上可以做两件事一是适当调大spark.executor.pyspark.memory默认值只有 512MB对大模型任务来说偏小二是控制单个分区的数据量避免一次性把整个分区的数据加载到内存。我在实践中通常把mapPartitions内部按 200 条一批处理每批处理完就释放内存占用稳定很多。另外输出结果尽量不要做“大字段的频繁转换”。比如把几万 Token 的长文本反复做 JSON 解析、正则替换、子串截取Python 进程的内存会很快膨胀。一次性加工完再返回比拆成多步操作更省内存。4.5 其他容易被问到的点最后说一个面试角度。Spark 怎么调用大模型这个问题面试官真正想听的其实不是requests.post怎么写在 UDF 里而是你能不能答出“为什么 UDF 直连不好”“并发怎么控制”“失败怎么处理”“结果怎么对齐”这几个关键点。我建议的答题结构是先说场景分类明确“这是离线批量推理”再说方案选型强调mapPartitions加客户端池化是工程上的主流选择最后补充并发、重试、成本统计这些落地的细节。如果面试官追问“为什么不能直接用 UDF”你就把连接复用、序列化、并发控制这三点讲清楚基本就能让人信服。如果面试官继续问“数据量上亿怎么办”你可以接上消息队列解耦和本地推理服务那套说明你确实从系统层面想过这个问题。这个问题现在几乎成了数据开发岗位的高频题准备好吃透固然重要但更关键的是你要有真实项目里的数据来判断方案的边界在哪里。5. 工程化落地监控、成本与稳定性治理5.1 可观测性日志、累加器与调用明细表代码能跑通只是第一步。真正放到生产环境里你还需要回答几个问题任务跑到哪了还有多少条没跑失败了多少条花了多少 Token如果要回答这些问题必须在设计的时候就埋好观测点。最简单的做法是用 Spark 的Accumulator来做进度统计。比如定义三个累加器分别统计成功条数、失败条数、重试次数from pyspark.sql import SparkSession spark SparkSession.builder.appName(spark_llm_inference).getOrCreate() success_counter spark.sparkContext.accumulator(0) fail_counter spark.sparkContext.accumulator(0) retry_counter spark.sparkContext.accumulator(0)然后在调用逻辑里根据结果更新累加器作业结束后通过success_counter.value就能看到最终的调用统计。这个方法比较简单但它的局限性也很明显累加器的值只能在作业结束后拿到任务跑的过程中你是看不到实时进度的。如果想做实时监控我建议把每次调用的明细写回一张结果日志表字段包括id、prompt长度、result长度、prompt_tokens、completion_tokens、latency_ms、status、error_msg、retry_count。这样既能看到单次调用的情况也能后续通过 SQL 做聚合分析。日志方面给一个实用建议不要在生产日志里打印完整的 prompt 和结果。一方面是因为数据敏感另一方面是日志文件会快速膨胀。我一般只打印id前缀、耗时、Token 数、状态码这些元信息。如果有特殊需要排查某一条数据可以通过id去查明细表。5.2 成本控制缓存、降级与模型分级调用大模型的成本是真实存在的尤其是公有云 API 按 Token 计费跑一次千万级数据的任务账单可能是几万甚至几十万。所以成本控制不是后端优化的问题而是方案设计阶段就要考虑的。第一个手段是输入缓存。很多真实业务场景里相似的文本会被重复提交。比如同一个文案在多个表里出现或者同一个用户的多条记录用的是同一段模板。可以在 Spark 作业开头先按文本内容做一次去重只对去重后的文本调用大模型再通过关联把结果映射回全量数据。如果业务场景里重复率高这个优化能把成本直接砍掉一半以上。第二个手段是结果缓存和增量处理。如果这个任务是每天跑的那不应该每次把全量数据都跑一遍。可以设计一个结果表记录每个业务 ID 最近一次处理的时间下次调度只处理新增和更新的数据。这样既节省费用也缩短了作业时间。第三个手段是模型分级。不是所有请求都需要最强的模型。简单的情感分类用 7B 模型就够了复杂的长文档摘要才需要 70B 甚至更大的模型。可以在数据准备阶段加一个路由逻辑按任务难度、输入长度、业务重要性把数据分成几档各调度到不同模型上。这个思路在业务量特别大的场景下很有价值也是我比较推荐的一种“架构级成本优化”。5.3 模型部署与迭代本地模型还是 API最后聊一下模型服务的形态。公有云 API 的好处是省心效果通常也比较好但问题是要考虑数据合规而且随着调用量增大成本是线性增长的。本地部署模型的好处是数据不出域、单 Token 成本低但运维负担重你既要做显存规划又要处理并发队列、负载均衡、模型版本的滚动升级。我的建议是先用公有云 API 跑通业务验证产品价值等技术方案稳定后再考虑本地部署。如果团队里有人熟悉 GPU 集群运维或者数据合规要求非常严格就提前做本地化的选型测试。现在用 VLLM 或者 TGI 部署开源模型已经非常简单接口也是 OpenAI 兼容的Spark 这边的代码几乎不用改只需要换一下base_url和model就行。这一点是这个方案最舒服的地方。如果你用的是 Spark 3.4 以上的版本并且 GPU 资源就在本地集群也可以尝试把 GPU 调度接入 Spark让 Executor 直接申请 GPU 资源跑推理。但对大多数团队来说把 GPU 推理做成独立的服务Spark 只作为客户端去调用是更简单、也更符合“大模型服务化”这个趋势的做法。后面模型要升级、要灰度、要回滚都在推理服务这一侧独立完成Spark 作业本身不会有感知。6. 最后一些个人体会文章写到这技术上的东西基本讲完了最后分享一点我自己的方法论。做这类项目我最大的体会是永远先跑一条链路再放大。不管方案设计得多完美第一次跑任务时一定先取几百条数据、用很小的分区数验证整个链路能不能通再逐步放大。我有一个同事就是跳过了这一步直接拿全量千万级数据跑结果 API Key 配错了作业跑了两小时才发现全是认证失败浪费了大量时间。另外一个体会是大模型调用的代码本身没那么难难的是把它嵌入到已有的数据链路里还不破坏稳定性。你要面对的是 Spark 的资源调度、Python 进程的内存边界、外部服务的限流和抖动、数据质量的参差不齐。这些问题的排查往往比写代码更花时间。所以我建议在做架构设计时一定要把监控、重试、限流、成本统计当成一等公民而不是最后补丁式地加上去。如果你正在做类似的项目希望这篇文章能帮你少踩几个坑。有什么更好的方案或者我没覆盖到的场景欢迎在评论区交流我这几年积累了不少稀奇古怪的报错样本说不定还能再写一篇问题排查合集。