
数据量一大几千上万条文本要逐个交给大模型处理又慢又费钱有没有更高效的路子这是最近一年我听到最多的问题。我自己长期在做数据平台和离线链路Spark是最常用的计算引擎恰好这段时间密集地接了大模型需求慢慢摸出一条在Spark任务里批量调用大模型的相对稳定路径。这篇把踩过的坑、对比过的方案、实际能跑的代码都整理出来给同样在数据管道里想接LLM的朋友一个可以直接参考的索引。它解决的核心问题是当你的数据存在Hive、HDFS、数仓表里量级在万级到亿级又需要逐条用大模型做分类、抽取、改写、扩写、打标时怎么用Spark的分布式能力把模型推理这件事并行化、批量化、稳定化并且尽量省钱省时间。适合数据工程师、算法工程师、以及所有想在离线链路里落地大模型能力的人。1. 为什么偏偏是 Spark 来接大模型1.1 数据规模决定了不能一条条手动处理很多团队的现状是数据已经在数仓里躺着了少则几万条多则几千万上亿条突然业务方说要给每条数据加一个AI摘要或者要做一遍敏感内容识别或者要按大模型的输出重新打标签。这时候如果写一个Python脚本在那里for循环单机跑一晚上只能处理几万条量一大直接崩。更现实的问题是你的数据本来就在Spark表里与其把几千万行导出去再导回来不如直接在Spark任务里把模型调用当作一个分布式计算步骤一条select就能把结果写回表。这种数据不动、计算走过去的思路是Spark来做模型调用的最大价值。1.2 离线批处理比在线接口更适合大模型场景大模型的单次响应通常在几百毫秒到几秒在线服务要求的是低延迟高并发而离线批任务天然不在乎单次要等多久在乎的是整体吞吐和成本。Spark这种把任务拆成很多partition每个任务并行跑的思路非常适合大模型这种IO密集加计算等待密集的场景。你开50个executor每个executor并行处理10条数据整体吞吐就上去了。配合Spark的Task重试机制某一条数据调用失败也不会拖垮整个任务。1.3 数据ETL与模型调用能在一个链路里完成过去做数据增强要先写Spark任务做预处理把结果导成JSON再用脚本调模型最后再把结果导回Hive中间至少裂成三段。现在直接在Spark里注册一个UDF把取数—清洗—调模型—解析结果—落表压成一条Spark SQL整条链路都在同一个计算引擎里完成血缘清晰、重跑方便。这个优势在需要频繁迭代提示词、频繁重算结果的项目里特别明显。2. 四类主流接入方案对比2.1 方案AUDF 外部API远程调用最直接的思路在Spark里写一个UDF每条数据构造一个HTTP请求发给线上大模型API。这个方案的优势是接入成本最低模型部署完全不用你操心调用阿里云、百度、智谱、DeepSeek等云厂商的OpenAI兼容接口就行。适合你的数据量不大十几万到百万级公司已经有统一的API网关或者你能接受每条数据几厘钱到几分钱的调用成本。但也有明显的坑网络IO会成为瓶颈。默认的HTTP客户端没有连接池每条数据新建连接慢不说还可能打到网关的限流上限。所以一旦决定走这个方案一定要做连接复用和重试绝不能裸写requests.post然后硬扛。2.2 方案BUDF 本地已部署模型服务如果你的团队自己部署了模型服务比如用vLLM、Ollama、TGI起了私有化接口那Spark这边的调用方式其实和方案A几乎一样只是endpoint从云厂商变成了内网地址。优势是数据不出内网合规性更好而且长文本、大批量时单位成本更低。缺点是模型服务本身需要GPU要做弹性伸缩否则高峰期会被Spark任务打爆。这里要提醒一句有些团队以为把Ollama装到Spark的executor所在机器上问题就解决了。实际不是。Ollama虽然能单机跑量化模型但多个executor同时把请求打到同一台机器的Ollama上显存和CPU都会瞬间拉满。比较稳妥的做法是模型服务独立部署在GPU机器上Spark只作为客户端。如果你本地是几块消费级显卡建议用Ollama起一个多模型常驻服务如果是企业级多并发优先上vLLM。2.3 方案CmapPartitions 多线程批量推理这是我在生产环境最推荐的方案也是性价比最高的。Spark里map算子是对每行执行一次mapPartitions是对每个partition执行一次。在mapPartitions里先把整个partition的数据收集成列表然后用线程池并发去请求模型服务等全部返回后再批量产出结果。这样做的好处是HTTP连接可以复用一次建连多次请求线程池并发度可以精准控制不会瞬间打死下游服务模型服务端到端的吞吐直接翻倍代价是实现复杂度上去了你需要自己处理线程池的创建、销毁、异常屏蔽、结果顺序映射。但对于几百万行以上数据这个改造能带来数倍到数十倍的性能提升。2.4 方案DGPU感知调度 分布式推理框架这是最硬核的方案简单说是让Spark任务直接调度GPU资源每个executor直接加载模型在GPU上推理或者通过Spark NLP这类框架把模型推理封装成分布式算子。优势是消除网络开销推理速度最快但工程复杂度也最高。你需要处理模型的分布式加载、显存分配、Executor动态变化导致的模型重新加载还要和YARN、Kubernetes做GPU资源调度集成。对于绝大多数团队来说性价比不如方案C因为你的GPU资源往往要留给训练任务让Spark任务直接吃GPU反而容易影响其他任务。四种方案对比如下维度方案A 外部API方案B 本地模型服务方案C mapPartitions批处理方案D GPU分布式推理部署难度低中低高性能中中高高最高稳定性受限于限流受限于服务容量可控受限于显存适用规模百万级以下百万到千万级千万级以上单批超大计算成本按量付费贵需维护GPU需维护GPU需深度集成推荐程度快速验证稳定通用生产主力特殊场景3. 最容易踩的坑序列化与资源管理3.1 模型对象为什么不能直接在UDF里传递很多人第一次写的时候都会踩这个我想在Driver端加载模型然后UDF里直接调用是不是就不用每次请求了基本原理不复杂Spark会把UDF的闭包序列化后分发到各个Executordriver里的Python对象默认不能直接复用。就算你用broadcast广播模型动辄几个GB广播变量的机制是把整个对象序列化后拉到每个executor内存里光序列化时间就够你睡一觉而且PySpark对象能否被pickle还是个问题。这里有一个我在生产里坚持的原则不要在Executor里加载大模型除非你用的是专门为分布式推理设计的框架。更合理的做法是让Executor只当客户模型统一放在独立的推理服务进程里Executor通过网络去请求。这样做的好处是模型的加载生命周期和Spark的executor生命周期解耦executor重启、扩缩容都不影响模型状态。3.2 Executor内存模型与OOM排查如果你试过在Executor里直接加载模型必然碰到过OOM。PySpark的Executor内存分两部分JVM堆内存spark.executor.memory和Python进程可用的内存。UDF里跑的是Python进程单独占用Executor节点上的系统内存默认并不受spark.executor.memory约束而是受spark.executor.pyspark.memory限制。这个参数一旦不设置Python端随便吃内存直到把整台机器打挂。我的建议是如果只做纯远程调用Python端几乎不占内存把pyspark.memory调小一些如果Executor端要缓存批数据就要估算单partition的数据量控制每个partition的行数别让线程池攒的数据全堵在一个executor里。平时排查OOM不能只看Spark UI的Storage内存要结合容器监控看Python进程的实际内存占用。3.3 YARN CPU分配异常是资源问题的常见征兆这里说一个大家问得特别多的现象我的Spark任务在YARN上跑CPU明明有几十核怎么每个executor只用1个vCore这多半不是模型调用的问题而是资源参数没配明白。Spark官方默认spark.executor.cores就是1也就是说每个executor只申请1个vCore即使机器再强也只用一核。你要做的是在提交时加--executor-cores 4或者调整spark.task.cpus让每个executor处理多个taskCPU才能跑起来。另外还有一种情况就是你发现executor数量很多但每个都很闲这不是配置问题而是并行度问题。数据量一定时分区数决定了任务数executor数乘以每个executor的CPU数才是真正的并行度。大模型调用场景下并行度也不是越大越好因为下游推理服务的并发能力是有限的Spark端开50个并发模型服务如果只能抗20个剩下30个全在排队白白消耗资源。所以在大模型场景下资源调优的重心其实不在Spark端而在Spark并发数这个模型服务吞吐量的匹配上。4. 完整实战从0到1实现Spark加模型批量推理4.1 环境准备与依赖引入我的生产环境是CDH 6.3.xSpark版本2.4/3.2都试过PySpark的用法基本一致。如果你用Spark 3.x配Python 3.8以上体验会更好。假设你现在有这样一个需求把Hive表里一列用户评论丢给大模型让它判断情感倾向是正/中/负同时抽取关键实体。本地已经用Ollama起了一个qwen2.5:14b模型服务地址是http://gpu-node-01:11434那么整个工程不需要额外引入什么重量级依赖只需要requests库用来发HTTP请求pyspark你的集群有就能用一个可以直接访问模型服务的网络环境如果你是调用OpenAI兼容的云厂商API把endpoint换成公网网关就行。我优先推荐先用Ollama验证等确定要走生产再切vLLM因为Ollama的接口足够简单排错也方便。4.2 第一版能跑的代码基础UDF调用本地模型下面这段代码是最基础的版本核心就是用PySpark注册一个UDF每条记录发一次HTTP请求到Ollama的/api/chat接口。我先用这种方式验证链路通不通。from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf from pyspark.sql.types import StringType import requests spark SparkSession.builder.appName(spark_llm_demo).enableHiveSupport().getOrCreate() MODEL_URL http://gpu-node-01:11434/api/chat def llm_sentiment(text): payload { model: qwen2.5:14b, messages: [{role: user, content: f判断下面评论情感 [{text}]只回答正向、中性、负向}], stream: False, options: {temperature: 0.2} } resp requests.post(MODEL_URL, jsonpayload, timeout120) resp.raise_for_status() data resp.json() return data[message][content] llm_udf udf(llm_sentiment, StringType()) df spark.sql(select id, comment from ads.user_comments where dt2024-12-01 limit 1000) result df.withColumn(sentiment, llm_udf(col(comment))) result.write.mode(overwrite).saveAsTable(tmp.user_comments_llm_result)这条链路跑通之后你会发现两个问题第一1000条数据可能要跑十几分钟因为每次请求都要新建HTTP连接而且串行执行第二如果有一批数据恰好碰到模型服务慢超时之后就整个任务失败。这就是基础版的局限所以它只适合用来验证连通性不适合上生产。4.3 第二版优化mapPartitions加线程池并发生产环境我建议直接上mapPartitions。下面这段代码是我在项目中反复调优后的一个模板from concurrent.futures import ThreadPoolExecutor, as_completed from pyspark.sql import SparkSession from pyspark.sql.functions import col import requests spark SparkSession.builder.appName(spark_llm_batch).enableHiveSupport().getOrCreate() MODEL_URL http://gpu-node-01:11434/api/chat def call_llm_once(item_id, text): payload { model: qwen2.5:14b, messages: [{role: user, content: f判断评论情感[{text}]只回答正向、中性、负向}], stream: False, options: {temperature: 0.2} } resp requests.post(MODEL_URL, jsonpayload, timeout120) resp.raise_for_status() content resp.json()[message][content] return item_id, content def batch_infer(iterator): rows list(iterator) if not rows: return futures [] with ThreadPoolExecutor(max_workers8) as pool: for row in rows: future pool.submit(call_llm_once, row[id], row[comment]) futures.append(future) result_dict {} for future in as_completed(futures): try: item_id, content future.result() result_dict[item_id] content except Exception as e: result_dict[item_id] fERROR:{str(e)[:100]} for row in rows: yield (row[id], result_dict.get(row[id], ERROR)) df spark.sql(select id, comment from ads.user_comments where dt2024-12-01) result df.mapPartitions(batch_infer).toDF([id, sentiment]) result.write.mode(overwrite).saveAsTable(tmp.user_comments_llm_result)这个版本的优化点有三个一是整个partition的数据只在开头收集一次网络请求复用同一个连接池requests库会自动keep-aliveSSL握手、TCP握手开销几乎消失二是8个线程同时发请求单条延迟被均摊一个partition的数据量越大均摊效果越明显三是单条失败不会让整个task崩溃异常会被截断成ERROR字符串保留下来方便事后分析。线程数的选择我是用探针法试出来的先设4看任务整体耗时和模型服务的GPU利用率然后逐步提高直到模型服务GPU利用率在80%左右就不再往上加了。在Ollama上qwen2.5:14b这种14B的量化模型8到12路并发是比较甜点的区间如果换成vLLM并发可以到30到50。4.4 第三版进阶让结果更稳定线上任务最怕的不是慢而是跑了一半某几个executor挂了。我建议再加三样东西第一是重试机制。单条调用失败后至少重试2到3次重试之间加一个小间隔能挡住很多临时性的网络抖动和服务繁忙。第二是写中间结果。如果数据量级到了千万以上建议不要等全部跑完再写表而是每隔一段时间把已完成的partition写成一个临时目录最后再合并这样任务中途被kill重启时能跳过已完成部分。第三是上下文压缩。如果你调用的是远程APItoken是成本大头可以把要传给模型的数据先做截断、去重再进Spark任务。下面这个伪代码展示了重试的写法import time def call_with_retry(text, retries3): last_exc None for i in range(retries): try: resp requests.post(MODEL_URL, jsonpayload, timeout120) resp.raise_for_status() return resp.json()[message][content] except Exception as e: last_exc e time.sleep(2 * (i 1)) return fERROR:{str(last_exc)[:100]}重试不是越多越好。重试次数超过3次大概率不是网络问题了而是模型服务真的不行这时候再重试只会把压力继续堆到服务端扩大雪崩。我在线上一般配置最多3次重试间隔采用指数退避。另外整个任务在spark-submit时要加spark.task.maxFailures默认4次如果某个任务的失败次数超过机器数可以调大一点给重试留足空间。4.5 关于Spark SQL的另一种优雅写法如果你不想写mapPartitions也可以把prompt构造放到SQL里用正则或concat函数先把输入拼好然后对生成的文本列注册一个只有一个参数的UDF。这样整条逻辑更贴合SQL驱动的团队习惯。示例from pyspark.sql.functions import concat, lit, col df spark.sql(select id, comment from ads.user_comments where dt2024-12-01) prompt_df df.withColumn(prompt, concat(lit(判断情感[), col(comment), lit(]只回答正向、中性、负向))) result prompt_df.withColumn(sentiment, llm_udf(col(prompt)))好处是提示词版本好管理出问题时可以在SQL里追溯坏处是单条并行度没有mapPartitions高因为每条数据仍然是一个UDF调用。所以我的建议是验证用UDF生产用mapPartitions。5. 常见问题与排查技巧实录5.1 spark on yarn提交是否需要单独安装Spark客户端这个问题在社区里反复被问到几乎每期必有人问。答案是只需要一个能做spark-submit的客户端节点就够了。YARN的模式下提交任务时客户端会把应用jar包和依赖上传到HDFS然后向ResourceManager发起申请ApplicationMaster会在某个NodeManager节点启动后续的计算任务全部由集群内部节点完成客户端提交完就可以退出不影响任务运行。真正要注意的是客户端环境要和集群版本匹配特别是Scala版本和Spark版本否则会出现各种NoSuchMethodError。5.2 为什么executor在yarn上运行时每个container只分配一个vCore这个问题的本质是你没有设置spark.executor.cores或设置了但没生效。默认值是1所以每个executor只申请1个vCore。在大模型调用场景下由于瓶颈往往在模型服务端Spark端的并行度设置其实可以保守一点。我推荐先设executor数量10个每个executor 2核也就是整体20并发观察模型服务的响应时间和任务完整体耗时然后慢慢调整。核心原则是不要让Spark端的并发数超过模型服务端可承受的并发数。5.3 模型冷启动慢导致第一批请求全部超时如果你用Ollama或者vLLM第一次请求时模型权重才加载进显存这个过程可能耗时几十秒到几分钟而当时Spark端的并发请求已经全部发出去很容易集体超时。我的处理办法是在Spark任务启动之前先发一条预热请求让模型服务把权重加载好。方法也很简单在spark-submit之前用curl访问一次模型接口curl http://gpu-node-01:11434/api/chat \ -d {model:qwen2.5:14b,messages:[{role:user,content:ping}],stream:false}预热请求返回之后再正式跑Spark任务。这个步骤能非常有效地避免跑了几分钟发现都是超时的尴尬。5.4 免费API限流怎么处理有些场景下你会用免费或低价的公共API这类API一般会限制每分钟调用次数。如果用了前面说的8线程并发很容易触发限流返回429或超时。我建议做两层控制一是mapPartitions里的ThreadPoolExecutor不要开太大4到6个就足够二是调用失败的请求要进入等待退避状态比如看到429就等30秒再重试。当然如果你愿意折腾也可以写一个基于Redis的分布式限流器让所有executor共享一个令牌桶这是更精细的做法但大多数场景用不上。把几个常见问题整理成速查表方便你排错现象可能原因解决办法任务一直PendingSpark并行度大于YARN可用vCore调小executor数量或cores某些executor一直GC模型返回结果过大缓存太多限制返回token数控制partition大小大量超时模型服务还在冷启动提前预热或增大UDF超时时间429错误触发了API限流降低线程并发加重试退避部分行结果是ERROR单条输入触发了模型安全策略单独排查这些行调整promptGPU利用率低但任务慢并发数不够调大线程池或executor数5.5 数据倾斜在模型调用场景怎么处理最后再说一个很多人忽略的问题数据倾斜。普通Spark任务数据倾斜会导致某个reduce卡很久模型调用场景也一样。假如有一类评论特别长模型推理时间可能是其他评论的10倍那么这个partition就会拖慢整个任务。我的处理办法是先把长文本单独拆出来走低并发的慢通道或者给文本长度加一列分桶按长度范围分别调度。还有一个更简单的思路是把超长文本先做摘要或者截断在保证效果的前提下大幅降低推理耗时。6. 一些亲测有效的配合技巧6.1 用异步调用进一步提升吞吐我在某些场景下还试过异步HTTP客户端比如用aiohttp在mapPartitions里把所有请求一次性并发发出不等待每个响应再统一收集。效果确实比线程池更高因为线程池在等待IO时依然会占着线程资源。不过异步代码的复杂度明显高很多对于大多数团队收益也没有那么夸张所以我把这个当成进阶技巧而不是默认方案。6.2 多模态模型接入时的额外考量如果你要处理的是图片、音频这类多模态数据调用的payload不是文本而是base64编码的文件内存占用会明显上升。我的建议是不要在Spark端把整个文件读成base64再传网络而是把文件路径传给模型服务让服务端自己读取。这样可以减少一半的内存和网络开销。前提是你的模型服务能访问到那个文件系统比如HDFS路径或对象存储路径。6.3 从成本维度倒推任务设计最后提醒一个大局观的问题。大模型推理是有成本的你用Spark每调一次API都在消耗token。因此在写Spark任务之前先想清楚这个问题这一列数据真的要全部调用大模型吗能不能先用规则过滤掉明显不需要处理的部分能不能只在增量数据上跑能不能用一个小模型先粗筛再把疑似样本交给大模型精判我见过太多项目几千万条数据全量跑一遍大模型费用高得吓人但其实80%的数据用正则就能处理。合理的任务设计永远比技术优化更省钱。我做这类任务的经验是先在本地小规模用基础UDF验证prompt效果确认输出格式稳定后再上mapPartitions批量方案最后再补重试和监控。先跑通再跑快最后跑稳。这套思路看着简单实际能帮你少踩很多坑。如果你正准备把大模型能力嵌入Spark管道希望这篇能帮你省下几个加班的夜晚。