
這個標題看起來像是一個純理論問題但背后是一個非常現實的工程需求團隊已經有成熟的 Airflow、Spark、dbt 之類的數據管道現在想在 ETL 里加一個“用 LLM 做文本分類、實體抽取、摘要、打標”的步驟。結果一接進去就發現問題延遲高、成本不穩、結果不可復現測試也不好寫。于是有人問能不能把 LLM 調用“編譯”成傳統數據管道那樣可重放、可緩存、可測試的確定性步驟。我的判斷是對于一批規則明確、輸入輸出穩定的 LLM 調用這條路是可行的而且有很清楚的工程范式。關鍵不是讓 LLM 變得更像數據庫而是把它包裝成“帶外部依賴的算子”然后套上緩存、批量、重試、回放這些數據管道本來就有機制。本文會把問題拆開給出一個最小可運行的管道示例再討論 API 接入、批量任務、資源占用、適用邊界和排查方法。1. 核心能力速覽在動手之前先把這個問題映射成工程參數。能力項說明問題類型LLM 應用架構設計不是具體開源工具核心問題平凡的 LLM 調用能否被改造成傳統數據管道的確定性步驟目標手段模板化、語義緩存、批量調用、失敗重試、確定性降級、蒸餾替換典型場景離線 ETL 中的文本分類、實體抽取、摘要、情感分析、標簽生成適合讀者有數據管道建設經驗想把 LLM 接入批處理的工程師不適合場景實時對話、Agent 多輪規劃、MCP 工具調用等強交互任務資源要求取決于本地推理還是 API 調用API 調用不占用本地 GPU批量能力支持按批次處理比逐條請求更容易控成本可測試性引入緩存和回放機制后可以大幅提升結果可復現性合規模塊數據脫敏、版權授權、隱私保護屬于必須前置條件這里要特別說明一點標題里的 “compiled” 不是傳統編譯器把高級語言變成機器碼的意思而是工程意義上的“把可變的外部調用改造成可驗證、可重放、可預測的數據處理步驟”。2. 先拆問題trivial LLM call 到底指什么“trivial LLM call” 翻譯過來是“很普通的 LLM 調用”。這種調用通常有幾個特征。第一輸入輸出結構非常固定。例如給定一條用戶評論輸出 positive/negative/neutral 三選一給定一段商品描述抽取品牌、型號、價格給一篇新聞輸出 100 字摘要。這類任務用提示詞模板就能做不需要復雜的多輪對話不需要工具調用也不依賴外部知識庫。第二業務邏輯相對穩定。今天是“評論情感三分類”下個月大概率還是這個分類體系。即使樣本變化任務目標不會頻繁變。這類任務很容易被固化成一個可重復執行的步驟。第三單次調用包含的信息量小。輸入可能只有一段文本輸出是一個短字符串或 JSON。這就是為什么在傳統管道里它會顯得很“廉價”也因此被稱為 trivial。但真正的問題是LLM API 畢竟是一次遠程調用它不是純函數。同樣一個 prompt溫度設為 0 也可能出現輸出抖動網絡超時會中斷整個管道費用隨調用量線性增長下游任務不知道這條數據是“新算出來的”還是“命中緩存復用的”。這些才是阻礙 LLM 調用進入傳統數據管道的真正原因。所以問題中 “compiled into conventional data pipelines” 的準確含義是能不能把這種普通 LLM 調用改造成數據管道里一個標準算子讓它擁有確定性、可重放性、可觀測性。3. 傳統數據管道為什么不接受原生 LLM 調用傳統數據管道里的核心組件從 SQL 的 UDF 到 Spark 的 Transformation到 Airflow 的 PythonOperator都有幾個默認前提結果可重放、運行成本可控、失敗時可重試、輸入輸出可結構化。LLM 原生調用在這些點上都不滿足。首先是不可復現。同一個輸入LLM 兩次調用的輸出可能不同。即使固定 temperature0不同模型版本、不同推理后端、甚至同一后端的量化參數都可能造成輸出漂移。數據管道下游往往需要穩定結果做聚合和對比輸出漂移會直接污染報表。其次是失敗模型不同。傳統管道失敗一般是數據缺失、類型錯誤、任務沖突LLM 調用失敗則是超時、限流、上下文超長、API key 失效、內容安全攔截。這要求管道具備完全不同的重試策略和降級策略。第三是成本不可預估。管道里處理 1 萬條數據和 1000 萬條數據SQL 的邊際成本幾乎為 0LLM API 的成本卻隨文本長度和調用量線性增長。如果管道沒有緩存和去重機制一次重跑就可能是賬單翻倍。第四是延遲。生產數據管道通常對吞吐有硬要求。如果中間插入一個逐條調用 LLM API 的算子整個管道的吞吐會瞬間被外部服務的響應延遲卡住。這四個問題合在一起結論已經比較明顯不是 LLM 不能進入數據管道而是需要給 LLM 套一層“編譯”機制先把上面四個問題解決掉。4. “編譯”在 LLM 管道里的三層含義把 LLM 調用編譯進傳統數據管道工程上一般分三層來做。這三層可以獨立實施也可以疊加使用。4.1 第一層模板化與語義緩存這是最快見效的一層。做法是把 LLM 調用封裝成一個算子輸入是一行結構化數據輸出是結構化字段中間只做一件事——用模板拼出 prompt然后調用 LLM最后解析輸出。同時給算子加磁盤緩存或語義緩存。緩存鍵可以是一整個 prompt 的哈希也可以是“任務名 輸入文本”的哈希。命中緩存就直接返回不發起 API 調用。這一步能把重復成本直接降到接近 0。對于數據管道來說同一條數據跑兩次、同一個批次被重放都是常見操作。如果沒有緩存每次重放都要為同樣的輸入付費。語義緩存比哈希緩存更進一步即使輸入文本有細微差異只要語義等價也能命中。但這個實現復雜度高需要向量化和相似度閾值適合在業務穩定后再引入。第一步先用精確匹配緩存性價比最高。4.2 第二層蒸餾與確定性替換這一層的思路更徹底如果某個 LLM 調用的任務是“穩定分類”或“穩定抽取”那就可以用第一批 LLM 輸出數據作為訓練集把它蒸餾成一個小模型、規則集甚至一段純 Python 代碼。這才是標題里 “compiled” 的最貼切含義——不是把 prompt 編譯成機器碼而是把一個“用自然語言描述的規則”編譯成確定性的代碼。舉個例子先用 LLM 批量標注 5000 條評論情感然后訓練一個很小的分類模型或者讓工程師從輸出中總結出一組關鍵詞規則。之后管道正式運行時就不需要再調 LLM直接跑小模型或規則即可。這一步適合高吞吐、低延遲、需要穩定輸出的場景。它的代價是前期需要一批標注數據和一個驗證流程。如果任務本身經常變蒸餾的成本可能比直接調用 LLM 還高。4.3 第三層編排系統里的 LLM 算子既不想完全蒸餾又想保留 LLM 的泛化能力那就把 LLM 調用封裝成管道里的標準算子并納入調度系統的重試、重放、監控體系。在 Airflow 里可以寫一個 PythonOperator內部批處理一批文本而不是逐條調用在 Spark 里可以用 mapPartitions 對每個分片批量調用在 Ray 里可以建一個遠程函數池并發調用。這一層解決的是“把 LLM 當普通計算單元”的問題。管道調度器不知道里面跑的是 SQL 還是 LLM API它只知道這個算子有輸入、有輸出、可能失敗、可以重試。5. 一個最小可運行的編譯型 LLM 管道示例紙上談兵沒有意義下面給一個可以直接跑通的最小示例。它演示了三個關鍵能力把 LLM 調用封裝成算子增加磁盤緩存避免重跑重復付費增加 fallback即使 LLM 調用失敗管道也不會整體中斷。生產環境把BaseLLMClient替換成真實 LLM API 客戶端即可。# llm_pipeline_demo.py import hashlib import json import time from dataclasses import dataclass from typing import Any, Callable, Dict, List, Optional class BaseLLMClient: 生產環境替換為真實 LLM API 客戶端。 def complete(self, prompt: str, temperature: float 0.0) - str: raise NotImplementedError class DummyLLMClient(BaseLLMClient): def complete(self, prompt: str, temperature: float 0.0) - str: time.sleep(0.05) # 模擬網絡延遲 return fresult_of: {prompt[:24]} class DiskCache: def __init__(self, cache_dir: str ./llm_cache): self.cache_dir cache_dir def _key(self, task: str, payload: Dict[str, Any]) - str: raw json.dumps( {task: task, payload: payload}, sort_keysTrue, ensure_asciiFalse, ) return hashlib.sha256(raw.encode(utf-8)).hexdigest() def get(self, task: str, payload: Dict[str, Any]) - Optional[str]: import os path os.path.join(self.cache_dir, self._key(task, payload) .json) if os.path.exists(path): with open(path, r, encodingutf-8) as f: return json.load(f)[output] return None def set(self, task: str, payload: Dict[str, Any], output: str) - None: import os os.makedirs(self.cache_dir, exist_okTrue) path os.path.join(self.cache_dir, self._key(task, payload) .json) with open(path, w, encodingutf-8) as f: json.dump({output: output}, f, ensure_asciiFalse) dataclass class LLMOperator: name: str llm: BaseLLMClient prompt_template: Callable[[Dict[str, Any]], str] cache: Optional[DiskCache] None fallback: Optional[Callable[[Dict[str, Any]], str]] None def execute(self, row: Dict[str, Any]) - Dict[str, Any]: # 1. 查緩存 if self.cache: cached self.cache.get(self.name, row) if cached is not None: row[self.name] cached row[f{self.name}_source] cache return row # 2. 構造 prompt 并調用 LLM prompt self.prompt_template(row) try: output self.llm.complete(prompt, temperature0.0) except Exception as exc: if self.fallback is None: raise output self.fallback(row) row[f{self.name}_error] str(exc) row[f{self.name}_source] fallback else: # 3. 寫入緩存 if self.cache: self.cache.set(self.name, row, output) row[f{self.name}_source] llm row[self.name] output return row def run_pipeline( rows: List[Dict[str, Any]], operators: List[LLMOperator], ) - List[Dict[str, Any]]: results [] for row in rows: current dict(row) for op in operators: current op.execute(current) results.append(current) return results def build_prompt(row: Dict[str, Any]) - str: return ( 請從以下評論中提取情緒positive/negative/neutral 和主題關鍵詞輸出 JSON。\n f評論{row[text]} ) def fallback_rule(row: Dict[str, Any]) - str: text row[text] if any(w in text for w in [好, 喜歡, 贊]): return {sentiment: positive} return {sentiment: neutral} if __name__ __main__: llm DummyLLMClient() cache DiskCache() op LLMOperator( namellm_extract, llmllm, prompt_templatebuild_prompt, cachecache, fallbackfallback_rule, ) rows [ {id: 1, text: 這個功能很好用我非常喜歡。}, {id: 2, text: 界面不太穩定經常卡頓。}, {id: 3, text: 功能正常速度可以接受。}, ] outputs run_pipeline(rows, [op]) for out in outputs: print(out)運行第二次所有llm_extract_source字段都會變成cache說明沒有再次發起 LLM 調用。這個模式非常簡單但它已經具備“編譯”的核心特征輸入確定、結果可緩存、失敗有降級、可以放進任何調度器。把DummyLLMClient換成真實客戶端就是這個思路的真實落地版本。接著可以把它接到 Airflow 風格的批處理任務里。以下是一個示意 DAG實際寫法需要按你的 Airflow 版本調整。from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def extract_with_llm(**context): batch_id context[params][batch_id] batch load_batch(batch_id) # 從數倉或文件讀取 results run_pipeline(batch, [llm_extract_operator]) write_to_parquet(results, foutput/batch_{batch_id}.parquet) with DAG( dag_idllm_compile_pipeline, start_datedatetime(2025, 1, 1), scheduledaily, catchupFalse, ) as dag: t1 PythonOperator( task_idrun_llm_batch, python_callableextract_with_llm, params{batch_id: {{ ds }}}, )注意這里最容易踩的坑是不要在 PythonOperator 內部逐條循環調用 LLM API。正確做法是每次執行處理一批數據內部再使用并發或分批請求。否則調度器一重試外部 API 會被同一批請求打爆。6. LLM API 接入、批量任務與成本控制真實項目不會用 DummyLLMClient下面講 API 接入時需要注意的工程細節。6.1 批量請求與隊列大多數 LLM API 都有限流策略單位時間內的請求次數和 token 數量都有限制。數據管道里最常見的錯誤是一個批次 1 萬條數據直接用 for 循環調用結果觸發限流任務中途失敗。推薦的做法是引入并發池加令牌桶。把每個批次切成小塊用concurrent.futures.ThreadPoolExecutor控制并發度。寫一個具備自動重試的通用調用函數比較穩妥。import openai # 按實際 SDK 安裝版本不同參數可能有差異 from tenacity import retry, stop_after_attempt, wait_exponential client openai.OpenAI() # 生產環境從配置或密鑰服務讀取 retry(stopstop_after_attempt(3), waitwait_exponential(min1, max10)) def chat_once(prompt: str) - str: resp client.chat.completions.create( modelgpt-4o-mini, # 按實際可用模型名替換 messages[{role: user, content: prompt}], temperature0.0, ) return resp.choices[0].message.content重試策略不要對超時、限流、網絡抖動和內容安全攔截一視同仁。限流通常需要等更長的時間內容安全攔截重試多少次都不會成功應該把它標記為錯誤數據落入異常隊列而不是無腦重試。6.2 緩存鍵與去重在數據管道里緩存鍵設計很重要。建議把“任務名 模型名 prompt 版本 輸入哈希”組合成緩存鍵。只對輸入文本做哈希容易在 prompt 模板升級后拿到舊結果。如果一次要處理 1000 萬條評論很多文本可能是重復的。先做一次去重再對唯一文本調用 LLM能用較少的請求覆蓋較大比例的數據。這個去重本身就是成本和速度的雙重優化。6.3 失敗重試與降級處理 LLM 調用失敗時至少要有三種策略。第一簡單重試。適合網絡抖動、瞬時限流。第二延遲退避。適合 API 側壓力大讓出時間窗口。第三確定性降級。如果業務允許失敗時用關鍵詞規則或默認值兜底。前面示例里的fallback_rule就是這個思路。降級結果必須標記來源否則下游會把它當成正常 LLM 輸出導致數據質量失真。完整的數據管道里LLM 調用的結論、來源、錯誤信息都應該作為列寫入結果表。這樣即使出現質量波動也能回溯到是哪一批數據、哪次 prompt 版本、哪個模型產生的問題。7. 資源占用與性能觀察資源占用要分兩種情況看。如果使用 LLM API本地管道資源主要是 CPU、內存和帶寬不占用 GPU。這時要重點觀察的是API 延遲的 P50/P95、失敗率、限流次數、緩存命中率、token 消耗量。建議每批次結束后寫一條監控日志至少包含這些指標。如果使用本地模型推理比如在管道里部署一個 7B 或 14B 的模型重點看顯存占用。顯存占用跟模型精度、batch size、輸入長度都有關系。使用 fp16/bf16 這類低精度格式能明顯降低顯存但模型效果可能需要在小樣本上驗證。這里不寫死某個模型具體占多少 GB因為不同量化等級、不同部署框架vLLM、llama.cpp、Ollama、Transformers差異很大需要以本機測試為準。觀察方法很簡單本地推理用nvidia-smi -l 1看實時顯存和 GPU 利用率API 調用在客戶端記日志管道層面在批次開頭和結尾記錄處理條數、耗時、失敗數。影響性能的主要變量有三個輸入文本長度越長token 消耗越大延遲越高并發度API 模式下并發越高吞吐越高但有限流邊界緩存命中率命中率越高平均單條成本越低吞吐越穩定。建議第一次跑通時先用 100 條數據做小批次性能測試觀察耗時和失敗率再逐步放大到全量。不要直接拿全量數據沖否則限流和成本都不可控。8. 適用場景與使用邊界這個思路適合什么場景適合輸入輸出結構固定、單次調用信息量小、業務規則相對穩定的批處理任務。比如用戶評論情感分類商品信息字段抽取新聞摘要客服工單打標郵件自動分類文檔版式識別后的內容清洗。不太適合什么場景不適合強交互、多輪決策、需要工具調用的復雜任務。Agent 規劃、MCP 工具調用、實時對話這類場景依賴狀態和多步反饋鏈路不可重放也缺少穩定的輸出結構硬塞進傳統批處理管道反而會把問題搞復雜。這類任務更適合用專門的 Agent 編排框架而不是把每次調用都“編譯”成固定算子。另外從 RAG 和語義檢索角度說如果管道里的 LLM 調用依賴外部向量庫或動態知識庫那輸入就不再是單行文本而是“文本 檢索上下文”。這種調用比 trivial call 復雜需要把檢索結果也納入緩存鍵和重放邏輯否則整個管道仍然不穩定。使用邊界還包括數據和合規。文本數據進入外部 LLM API 前必須確認數據是否包含個人隱私、商業秘密、版權內容。涉及人臉、聲音、肖像、受版權保護的素材時必須確認授權。生產環境建議先做脫敏再評估能否使用外部 API。如果數據不能出域就要選擇本地部署推理服務成本和治理都不同。9. 常見問題與排查方法問題現象可能原因排查方式解決方案管道跑完發現很多結果來自 fallbackLLM 調用失敗但被降級處理檢查結果表中的_source和_error字段區分“LLM 正常結果”和“降級結果”統計失敗率同一份數據重跑成本翻倍緩存鍵設計不合理或緩存未生效檢查緩存命中率日志把任務名、模型名、prompt 版本納入緩存鍵API 請求大量超時并發度過高或限流查看客戶端日志和 API 錯誤碼降低并發度增加指數退避重試結果不穩定影響下游報表輸出解析失敗或模型輸出抖動對比同一條輸入的多次輸出固定 temperature0增加輸出格式校驗和重試顯存不足或推理很慢本地模型精度或 batch size 設置不當用nvidia-smi觀察顯存記錄單批次耗時換低精度加載減小 batch size或改用 API模型升級后結果風格變化prompt 模板或模型版本變更對比歷史輸出每次升級前用固定測試集做回歸對比大批量處理時中間斷掉沒有做批次檢查點和斷點續跑檢查調度器日志按批次寫結果重跑時跳過已完成批次輸出 JSON 解析失敗模型輸出包含多余文本記錄原始輸出增加輸出格式約束或解析后校驗失敗重試10. 最佳實踐與使用建議先給一個保守的落地路徑。第一步把一個 LLM 調用封裝成算子加入緩存和降級用小批量數據驗證正確性。不要先上完整管道先驗證單算子。第二步把算子接入現有調度器但保留“跳過 LLM 直接跑緩存”的模式。這樣在 prompt 調整和模型升級時能對比新舊結果。第三步建立固定的評測集。至少準備 100 到 500 條帶標準答案的樣本每次模型或 prompt 變更后都跑一遍對比準確率和格式合格率。沒有評測集的 LLM 管道后期維護會非常痛苦。第四步做蒸餾或規則替換。當批量任務穩定運行一段時間后把高頻輸入和 LLM 輸出導出嘗試用規則、小模型替換。目標是把 80% 的確定性請求從 LLM 調用中剝離出去只保留少數復雜樣本走 LLM。第五步監控成本和質量。每批次記錄 token 消耗、緩存命中率、fallback 次數、輸出解析失敗率。成本和質量一旦異常能快速定位是數據變化、prompt 變化還是模型變化導致。關于 LLM 文本向量 API 未配置這類問題如果管道里還涉及向量化、語義檢索建議把“LLM 調用”和“向量化調用”分開配置、分開監控。混在一起會導致故障定位困難尤其是其中一方限流或 key 失效時很難判斷是哪個服務導致整個管道卡住。11. 總結與下一步“Can trivial LLM calls be compiled into conventional data pipelines?” 這個問題的答案不是簡單的“能”或“不能”。對于輸入輸出固定、業務規則穩定的普通 LLM 調用答案是“能但需要經過工程改造”。改造的核心是把 LLM 從“隨時可能抖動的外部服務”封裝成“帶緩存、帶重試、帶降級的管道算子”。最先應該驗證的功能是緩存能否正確命中、降級路徑能否在不中斷管道的情況下兜底。最容易踩的坑是直接逐條調用 API以及不記錄結果來源導致下游無法判斷數據質量。后續可以從已有管道里挑一個最簡單的文本分類任務開始先跑通單算子再做批次調度和成本監控。如果這個方向驗證成功下一步可以把實驗擴展到模板化生成、摘要、字段抽取如果業務穩定到一定程度再考慮把 LLM 輸出蒸餾成確定性組件徹底擺脫對在線模型的依賴。先跑通最小示例再逐步把評測集、影子比較、成本監控加進去是比較務實的路徑。