
最近在 Hacker News 上看到一个讨论Can trivial LLM calls be compiled into conventional data pipelines翻译过来就是——那些非常简单的 LLM 调用能不能像普通代码一样被“编译”进传统的数据管道里。这个问题的切入点很务实。过去一年里很多团队把 LLM 当成一个“高级文本函数”来用给一段文本拿回一段 JSON。做法通常是写一个脚本循环调用接口把结果写成 CSV 或 JSONL。功能跑通了但维护起来很吃力——没有缓存、没有重试、没有校验跑完也说不清到底花了多少 token、哪些记录失败了。本文会从概念、设计到完整代码把这个话题拆开讲清楚。先解释什么样的 LLM 调用属于 trivial 型为什么把这类调用管道化有价值再带你实现一个可运行的示例读取 JSONL 评论数据用 LLM 抽取结构化字段经过校验后写回文件。代码基于 Python 和 OpenAI 兼容接口换成 vLLM、Ollama 或其他兼容服务同样适用。1. 概念先行什么样的 LLM 调用才算“平凡”1.1 trivial LLM call 是什么“trivial”在这里指“平凡、简单、不需要复杂推理”的调用。这类调用通常具备以下几个特征输入是一段文本输出是一个可解析的文本结果。不需要多轮对话不需要维护会话状态。不依赖外部工具、数据库查询或额外的记忆。所有必要信息都包含在 prompt 里模型只需要“读进来、写出去”。典型场景包括从文本里抽取关键词和字段、判断情感倾向、翻译润色、把非结构化文本转换成固定 JSON、给评论打标分类。比如输入“相机画质很好但电池续航一般”希望输出{sentiment: positive, category: 相机, keywords: [画质, 电池续航]}这就是一个很典型的 trivial LLM call。从数据工程视角看这类调用本质上和“把字符串转成数字”没有区别只是函数内部变成了一个模型。它不涉及 Agent 那种“边思考边行动”的复杂循环也不需要在推理过程中不断调整 prompt。正因为逻辑简单它才有机会被抽象成管道里的一个固定算子。1.2 传统数据管道有哪些特征传统数据管道通常描述这样一类系统从数据源拉取数据经过清洗、转换、校验最终写入目标存储。无论你用 Airflow、dbt、Flyte 还是简单的 cron Python 脚本管道都应该具备几个基本特征分阶段每个阶段职责单一比如抽取、转换、加载。幂等同一份数据重复跑一遍不会产生重复结果。可重试中间环节失败可以安全地从断点继续。可观测能记录处理数量、耗时、失败原因。可测试输入相同输出可预期。对比一下常见的 LLM 调用脚本会发现它和管道思维有天然的距离。普通脚本往往是串行循环逐条调用接口结果直接落盘没有中间状态也没有失败分区。对于几十条数据的临时任务这样写没问题对于每天几万条、几十万条的生产任务问题就会集中爆发。下面用一个表格直观对比两种方式的差异维度一次性调用脚本管道化处理缓存无相同请求重复调用按 prompt 模型做缓存重试手动处理或没有指数退避自动重试校验无输出 JSON Schema 校验可观测性无记录 token、延迟、成功率幂等性无按原始 ID 去重写入并发能力串行速度慢线程池或批处理并行1.3 为什么要“编译”成管道“编译”这个词在这里是比喻不是说把 LLM 调用变成机器码而是把“运行时临时发起的调用”变成一种确定性的、可复用的管道阶段。这样做的主要原因有三个。第一是成本控制。LLM 调用按 token 计费同一个 prompt 如果被重复调用就是在重复烧钱。缓存命中一次省下的费用可能很可观。第二是稳定性。接口会限流、会超时、会偶尔返回非法 JSON这些都需要管道层面的兜底而不是让脚本直接崩溃。第三是可维护性。把 prompt 模板、输出 schema、校验规则、模型版本固化下来数据工程师才能像维护普通 ETL 一样维护 LLM 处理流程。当然也不是所有 LLM 调用都适合管道化。多轮 Agent、动态工具调用、需要流式输出的实时对话、依赖外部检索的 RAG 场景强行套管道会损失灵活性。本文后面会专门讨论适用边界。2. 环境准备与项目结构为了让示例能直接运行建议使用以下环境Python 3.10 及以上版本。openai Python SDK 1.x用于调用 OpenAI 兼容接口。SQLite3Python 内置用于实现本地响应缓存。可选 tiktoken用于估算 token 用量。需要说明的是版本号不需要严格锁定。如果你使用的是国内大模型服务或私有化部署的 vLLM、Ollama只需把 base_url 和 api_key 换成对应配置即可。本示例的核心是管道设计思路而不是某个封闭框架。项目结构如下llm_pipeline_demo/ ├── data/ │ ├── input.jsonl │ └── output.jsonl ├── llm_client.py ├── run_pipeline.py └── requirements.txt依赖文件内容openai1.0 tiktoken0.7安装命令pip install -r requirements.txt如果你使用本地推理服务安装完成后请确认模型服务已经启动并准备好兼容的 API 地址。后续所有代码都通过OpenAI()客户端调用只要接口协议兼容逻辑完全一致。3. 核心设计把 LLM 调用变成管道算子3.1 输入输出契约要让一段 LLM 调用成为管道算子首先要定义好“输入是什么输出是什么”。没有契约函数和函数之间就没法组合。在示例中输入一条评论记录{id: 1, text: 相机画质很好但电池续航一般。}输出一条带抽取结果的新记录{ id: 1, text: 相机画质很好但电池续航一般。, extraction: { sentiment: positive, category: 相机, keywords: [画质, 电池续航] }, valid: true }这里的关键是把 prompt 也视为代码的一部分。prompt 模板、temperature、模型名称、输出 JSON 结构都属于这个算子的“编译参数”。后续无论谁来调用都必须使用同一套参数管道行为才可预期。一个标准的 prompt 模板可以写成你是商品评论数据清洗助手。 请从评论中抽取 - sentimentpositive / negative / neutral - category商品品类 - keywords关键词数组 只输出合法 JSON格式为 {sentiment: ..., category: ..., keywords: []} 评论内容 {text}注意这个模板里明确规定了输出格式和字段值域这比只写一句“抽取关键信息”要可靠得多。LLM 任务里prompt 的约束力直接决定结果的稳定性。3.2 从“即时调用”到“编译后阶段”把上面的模板、模型参数、缓存策略、校验逻辑组合起来就得到了一个可复用的管道算子。它的执行过程可以抽象成下面这条链路输入记录 → 构造 messages → 查缓存 → 命中则直接返回 ↓ 未命中 调用模型 ↓ JSON 解析与校验 ↓ 写入缓存 → 返回结果对比原始的一次性脚本这里多出来的几层都很关键。缓存层负责消除重复成本重试层负责应付限流和超时校验层负责拦截模型的脏输出。三层叠加之后LLM 调用从一个“不可控的黑盒”变成了一个“边界清晰的处理单元”。需要提醒的是这里的“编译”是概念层面的抽象。实际编码时我们仍然是用 Python 函数来表达这个算子只是它具备了可缓存、可重试、可校验的特性。对于数据工程师来说这已经足够把它当作普通转换函数嵌进更大的管道。3.3 适用边界在动手写代码之前先明确哪些任务适合这种处理方式。场景是否适合管道化原因批量抽取字段、打标、分类非常适合逻辑简单输入输出固定翻译、改写、润色适合单条文本独立处理短文本摘要适合不依赖外部检索多轮 Agent 对话不适合状态和分支太多动态工具调用不适合依赖执行结果再决策实时流式聊天不适合延迟敏感不适合批处理RAG 问答谨慎需要检索和拼装上下文已超出“平凡调用”范围判断标准其实很简单如果一次调用在发起前就确定了完整的 prompt并且输出不需要再回传给模型继续推理那它就是平凡调用适合管道化。如果中间需要根据结果决定下一步动作那它已经进入 Agent 领域应该单独设计编排层。4. 完整实战评论数据抽取管道下面我们搭建一个最小可运行的示例。它会把data/input.jsonl里的评论逐条读取交给 LLM 抽取结构化字段校验后写入data/output.jsonl同时用 SQLite 缓存相同请求。4.1 准备示例数据创建data/input.jsonl{id: 1, text: 相机画质很好但电池续航一般。} {id: 2, text: 物流很快客服态度也不错整体满意。} {id: 3, text: 收到就坏了差评。} {id: 4, text: }第 4 条故意留空用来验证管道对异常输入的处理。4.2 实现带缓存的 LLM 客户端创建llm_client.py# llm_client.py import json import hashlib import sqlite3 import time import logging from openai import OpenAI logger logging.getLogger(__name__) # 读取环境变量 OPENAI_API_KEY也可以显式传入 base_url 和 api_key client OpenAI() def _get_conn(db_pathllm_cache.db): conn sqlite3.connect(db_path) conn.execute( CREATE TABLE IF NOT EXISTS llm_cache ( cache_key TEXT PRIMARY KEY, response TEXT, model TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ) ) return conn def make_cache_key(messages, model, temperature): payload json.dumps( {messages: messages, model: model, temperature: temperature}, ensure_asciiFalse, sort_keysTrue, ) return hashlib.sha256(payload.encode(utf-8)).hexdigest() def call_llm_with_cache(messages, modelgpt-4o-mini, temperature0.2, max_retries3, db_pathllm_cache.db): key make_cache_key(messages, model, temperature) conn _get_conn(db_path) try: row conn.execute( SELECT response FROM llm_cache WHERE cache_key ?, (key,) ).fetchone() if row: logger.info(cache hit, key%s, key[:12]) return json.loads(row[0]) for attempt in range(1, max_retries 1): try: resp client.chat.completions.create( modelmodel, temperaturetemperature, response_format{type: json_object}, messagesmessages, ) content resp.choices[0].message.content result json.loads(content) conn.execute( INSERT OR REPLACE INTO llm_cache(cache_key, response, model) VALUES (?, ?, ?), (key, content, model), ) conn.commit() return result except json.JSONDecodeError as exc: logger.warning(JSON parse error, attempt%s, error%s, attempt, exc) if attempt max_retries: raise except Exception as exc: logger.warning(LLM call failed, attempt%s, error%s, attempt, exc) if attempt max_retries: raise time.sleep(min(2 ** attempt, 30)) return None finally: conn.close()这里有几个设计点需要解释。make_cache_key对 messages、model、temperature 做序列化哈希只要这三者完全一致就能命中缓存。序列化时使用sort_keysTrue保证字段顺序不影响 key 值。缓存表用 cache_key 作为主键重复写入时用INSERT OR REPLACE覆盖避免脏数据残留。重试逻辑采用指数退避第一次失败等 2 秒第二次 4 秒最多等 30 秒比固定间隔更适合限流场景。如果你的模型服务不支持response_format参数可以删掉这一行但必须通过更严格的 prompt 和后处理来保证 JSON 输出否则解析失败率会明显上升。4.3 实现校验与主流程创建run_pipeline.py# run_pipeline.py import json import argparse import logging from concurrent.futures import ThreadPoolExecutor, as_completed from llm_client import call_llm_with_cache logging.basicConfig(levellogging.INFO, format%(asctime)s %(levelname)s %(message)s) logger logging.getLogger(__name__) SYSTEM_PROMPT ( 你是商品评论数据清洗助手。 请从评论中抽取 sentimentpositive/negative/neutral、 category商品品类、keywords关键词数组。 只输出合法 JSON不要输出任何解释。 ) USER_TEMPLATE 评论内容{text} def build_messages(text): return [ {role: system, content: SYSTEM_PROMPT}, {role: user, content: USER_TEMPLATE.format(texttext)}, ] def validate_extraction(data): if not isinstance(data, dict): return False if data.get(sentiment) not in {positive, negative, neutral}: return False if not isinstance(data.get(keywords), list): return False return True def process_one(record): record_id record.get(id) text record.get(text, ).strip() if not text: return {id: record_id, valid: False, error: empty_text} messages build_messages(text) try: extraction call_llm_with_cache(messages) valid validate_extraction(extraction) return { id: record_id, text: text, extraction: extraction if valid else None, valid: valid, } except Exception as exc: return {id: record_id, text: text, valid: False, error: str(exc)} def main(): parser argparse.ArgumentParser(descriptionLLM 数据抽取管道) parser.add_argument(--input, defaultdata/input.jsonl) parser.add_argument(--output, defaultdata/output.jsonl) parser.add_argument(--workers, typeint, default4) args parser.parse_args() records [] with open(args.input, r, encodingutf-8) as f: for line in f: line line.strip() if line: records.append(json.loads(line)) results [] with ThreadPoolExecutor(max_workersargs.workers) as executor: futures {executor.submit(process_one, r): r for r in records} for future in as_completed(futures): result future.result() results.append(result) logger.info(processed id%s valid%s, result.get(id), result.get(valid)) results.sort(keylambda x: x.get(id, 0)) with open(args.output, w, encodingutf-8) as f: for result in results: f.write(json.dumps(result, ensure_asciiFalse) \n) valid_count sum(1 for r in results if r.get(valid)) logger.info(done, total%s valid%s invalid%s, len(results), valid_count, len(results) - valid_count) if __name__ __main__: main()主流程的结构很简单先读取输入文件然后用线程池并发处理最后按 id 排序写回输出文件。处理每条记录时先把空文本拦截掉再构造 messages、调用带缓存的 LLM 客户端、执行 schema 校验。任一步失败都会被捕获并记录到error字段不会影响其他记录继续处理。并发写入 SQLite 在示例规模下没有问题但如果数据量很大建议把缓存单独放进 Redis或者把 SQLite 连接设置为单写模式这部分在第 5 节会细讲。4.4 运行与验证先设置 API Keyexport OPENAI_API_KEYsk-你的key然后运行管道python run_pipeline.py --input data/input.jsonl --output data/output.jsonl --workers 4第一次运行预期日志类似2025-01-10 12:00:01 INFO processed id1 validTrue 2025-01-10 12:00:02 INFO processed id2 validTrue 2025-01-10 12:00:03 INFO processed id3 validTrue 2025-01-10 12:00:03 INFO processed id4 validFalse 2025-01-10 12:00:03 INFO done, total4 valid3 invalid1第二次运行同样的命令你会看到cache hit日志说明缓存生效不再产生新的模型调用。输出文件data/output.jsonl大致如下{id: 1, text: 相机画质很好但电池续航一般。, extraction: {sentiment: positive, category: 相机, keywords: [画质, 电池续航]}, valid: true} {id: 2, text: 物流很快客服态度也不错整体满意。, extraction: {sentiment: positive, category: 物流, keywords: [物流, 客服]}, valid: true} {id: 3, text: 收到就坏了差评。, extraction: {sentiment: negative, category: 其他, keywords: [坏]}, valid: true} {id: 4, text: , valid: false, error: empty_text}4.5 结果说明这个示例虽然简单但已经具备了管道应有的几个基本能力空数据被提前拦截模型输出被 schema 校验失败记录不会让整个任务崩溃重复运行能命中缓存。它还展示了并发处理的写法让多条记录可以并行调用模型而不是一条条排队。如果想继续扩展最直接的改动是把输出从 JSONL 换成数据库表或者把process_one里的逻辑拆成独立的抽取、校验、加载三个子函数这样每层都能单独测试。也可以把input.jsonl换成消息队列让管道变成消费式的流处理任务。5. 常见问题与排查思路LLM 管道和普通 ETL 最大的区别在于它的“计算单元”是一个外部网络服务不确定性更高。下面汇总几个高频问题。问题现象常见原因解决思路429 rate limit并发过高或配额不足指数退避重试、降低并发数、离线任务使用 Batch APIJSON 解析失败模型输出 markdown 或多余解释使用 response_format、加强 system prompt、增加修复重试缓存命中率低prompt 拼入无关动态字段或顺序变化固定字段顺序、归一化输入、忽略无关空白输出字段漂移模型版本更新或 prompt 调整固定模型版本、记录 prompt version、加 schema 校验成本超预期重复数据处理、未做缓存增加缓存、增量处理、先用小模型分流SQLite database is locked多线程并发写同库开启 WAL、单写多读、或改用 Redis5.1 限流与 429限流是调用外部模型最常遇到的问题。报错长这样openai.RateLimitError: Error code: 429排查顺序是先看是不是并发数设得过高再看是不是配额不足。示例代码里的指数退避可以缓解瞬时限流但如果每分钟调用量超过账号上限就需要降低--workers或者把离线批量任务切到 Batch API这样不仅更便宜限流压力也小很多。5.2 JSON 解析失败即使要求“只输出 JSON”模型也有概率输出带 json 包裹的文本。稳健的做法是先用response_format开启 JSON 模式再在解析失败时做一次兜底清洗。def robust_json_loads(content): content content.strip() if content.startswith(): content content.strip() if content.startswith(json): content content[4:] return json.loads(content)如果清洗后仍然解析失败可以把原始响应和报错信息拼接回 messages让模型“参考上面的错误重新输出一次”。这个修复重试机制在生产环境中非常实用。5.3 缓存命中率低缓存 key 必须稳定才能保证命中率。常见的坑包括prompt 里带入了当前时间、随机数、或字典序不固定的字段。示例中使用sort_keysTrue并且只对 messages、model、temperature 做哈希就是在刻意减少不稳定因素。另一个坑是输入文本的空白差异。如果同一句话一会儿带换行、一会儿不带缓存 key 就会不同。建议在进入管道前统一做文本清洗比如去掉首尾空白、统一空格。5.4 输出字段漂移模型版本更新后同样的 prompt 可能返回不同的字段值。解决方案有三层第一固定模型版本比如写死gpt-4o-mini-2024-07-18这类完整版本号第二在输出记录里写入 model 和 prompt_version方便追溯第三用 JSON Schema 做严格校验不满足就标记为无效记录而不是直接落库。5.5 SQLite 并发写锁多线程并发写 SQLite 时可能出现database is locked。简单办法是给连接开启 WAL 模式conn.execute(PRAGMA journal_modeWAL)WAL 模式允许多个读并发和一个写并发适合示例这种轻量场景。如果数据量再往上走建议把缓存迁到 Redis或者把写库操作收敛到单一线程。6. 最佳实践与工程建议6.1 把 prompt 和 schema 纳入版本管理prompt 是 LLM 管道里最容易被随意修改的部分。建议把 prompt 模板、输出 schema、模型版本统一放在一个配置模块或文件里并给每次改动增加prompt_version字段。输出记录里带上版本号将来排查“为什么结果变了”会轻松很多。6.2 缓存设计要考虑成本与一致性缓存 key 必须包含 messages、model、temperature但不要包含与结果无关的元信息。如果业务允许还可以在缓存层之前做文本相似度去重避免两条几乎相同的文本重复调模型。需要注意的是一旦 prompt 模板或模型版本升级旧缓存应该自动失效否则线上可能一直读到旧结果。6.3 可观测性要落实到记录级别普通脚本只需要关心“成没成功”LLM 管道还需要关心“花了多少钱、多慢、哪条失败了”。建议在每一条记录处理完成后输出结构化日志包含记录 ID、模型、prompt token、completion token、延迟、是否命中缓存、是否校验通过。usage resp.usage logger.info( tokens prompt%s completion%s total%s, usage.prompt_tokens, usage.completion_tokens, usage.total_tokens, )有了这些日志就能在一天跑批结束后快速回答“这次处理花了多少钱”“哪些模型调用最耗时”“失败集中在哪个数据源”这类问题。6.4 数据安全边界调用外部模型时必须明确数据边界。如果输入包含手机号、地址、身份证等敏感信息不要直接拼进 prompt。流程上可以先做脱敏或者使用私有化部署的模型服务。API Key 一律通过环境变量或密钥管理服务注入不要写死在代码里。对模型输出也要做二次内容安全过滤因为模型结果不一定符合业务合规要求。6.5 生产环境的上线顺序上线 LLM 管道前建议按以下顺序检查先用几十条典型样本离线验证 prompt 和 schema估算每天的数据量、token 消耗和成本设计好增量处理逻辑记录已处理过的记录 ID把失败记录写入死信队列或单独文件方便人工复查最后再加上监控告警关注成功率和耗时。不要一上来就全量重跑历史数据。先从增量数据开始观察几天稳定后再回补历史这是更稳妥的做法。7. 总结与学习路线这篇文章的核心观点是trivial LLM 调用在逻辑上非常接近普通转换函数完全可以被“编译”成传统数据管道中的一个算子。关键在于给它定义好输入输出契约并补上缓存、重试、校验、观测这几层工程能力。示例中的评论抽取管道虽然简单但已经展示了从 JSONL 读取、并发调用、SQLite 缓存、schema 校验到结果落盘的全过程。如果你准备在真实项目中落地下一步建议优先了解这些方向数据加载工具 dlt 已经支持把 LLM 调用声明式地嵌入管道LangChain 和 LlamaIndex 提供了更丰富的语义编排能力离线大批量处理可以看看 OpenAI Batch API 或 vLLM 的离线推理调度层面则可以用 Airflow、Flyte、Kubeflow 这类成熟平台。重点还是要区分好“平凡调用”和“Agent 任务”前者适合管道化后者需要独立设计状态和分支逻辑。建议你直接把示例代码跑通再尝试改成你自己的抽取任务比如从工单里提取优先级、从日志里识别故障类型、从商品描述里抽取规格参数。跑通之后再逐步加入增量处理、监控和成本统计你会明显感受到 LLM 数据管道化带来的稳定性提升。如果这篇文章对你有帮助可以收藏备用也欢迎在评论区分享你在 LLM 数据处理上的踩坑经验。