Clojure 集成大语言模型

Clojure 集成大语言模型的完整工程实践:多供应商 API 封装与 SSE 流式解析、core.async 管道编排与背压、工具调用与结构化输出校验、向量检索与 RAG 拼装、token 计量与并发限流、两级缓存去重,以及可观测指标与评测回放方法。

LLM 应用的技术栈里,Clojure 不是最显眼的选择,但它有几个被低估的优势:数据导向让 prompt 与响应天然是 map,core.async 天生适合编排流式与并发,不可变数据让重放与评测变得简单。缺的是工程化的封装——API 客户端、流式解析、工具调用、RAG、成本控制,这些都得自己搭。本文按「封装 → 流式 → 工具 → 检索 → 成本 → 可观测」的顺序,给出一套可运行的 Clojure LLM 集成方案。

1. Clojure 做 LLM 编排的定位

1.1 适合做什么

场景契合度原因
多供应商路由与降级高数据驱动配置 + 纯函数转换
流式管道编排高core.async 天然背压与并发
RAG 数据加工高数据处理是 Clojure 强项
评测与回放高不可变数据 + 纯函数
模型训练/微调低那是 Python 的领域

定位建议:Clojure 做编排层与数据层,把重计算(推理、微调)交给专门的服务。这样既能吃到 Clojure 的数据处理优势,又不必跟生态硬碰。

1.2 依赖选型

;; deps.edn
{:deps {org.clojure/clojure            {:mvn/version "1.12.0"}
        org.clojure/core.async         {:mvn/version "1.6.681"}
        metosin/malli                  {:mvn/version "0.16.4"}
        http-kit/http-kit              {:mvn/version "2.8.0"}
        cheshire/cheshire              {:mvn/version "5.13.0"}
        com.github.seancorfield/next.jdbc {:mvn/version "1.3.939"}}}

HTTP 客户端用 http-kit(异步、支持流式),JSON 用 cheshire,schema 用 Malli。这些和 Clojure 数据管道与流处理 里的技术栈一致,可以复用同一套基础设施。

2. API 封装

2.1 统一客户端抽象

不要在每个业务函数里直接拼 HTTP 请求。先定义一个供应商无关的接口:

(ns app.llm.protocol)

(defprotocol LLMClient
  (complete [this request])        ; 一次性返回
  (stream   [this request])        ; 返回 core.async 通道
  (embed    [this texts]))         ; 向量化

request 是纯数据:

{:model    "gpt-4o-mini"
 :messages [{:role :system :content "你是助手"}
            {:role :user   :content "..."}]
 :tools    [...]        ; 可选
 :temperature 0.2
 :max-tokens 1024}

好处:换供应商只改实现,业务代码不动;测试时用 mock 实现即可,无需网络。

2.2 OpenAI 兼容实现

(ns app.llm.openai
  (:require [app.llm.protocol :as p]
            [cheshire.core :as json]
            [org.httpkit.client :as http]))

(defn- ->payload [{:keys [model messages tools temperature max-tokens]}]
  (cond-> {:model model :messages messages}
    tools        (assoc :tools tools)
    temperature  (assoc :temperature temperature)
    max-tokens   (assoc :max_tokens max-tokens)))

(defn complete [_ req]
  (let [{:keys [body status]}
        @(http/post "https://api.openai.com/v1/chat/completions"
                    {:headers {"Authorization" (str "Bearer " (api-key))
                               "Content-Type" "application/json"}
                     :body (json/generate-string (->payload req))
                     :timeout 60000})]
    (if (= 200 status)
      (-> (json/parse-string body true)
          (get-in [:choices 0 :message]))
      (throw (ex-info "llm request failed"
                      {:status status :body body})))))

2.3 错误、重试与超时

LLM API 的错误分三类,处理方式完全不同:

类型HTTP处理
限流429指数退避重试,尊重 Retry-After
服务端错误5xx退避重试,可降级到备用模型
请求错误4xx不要重试,改 prompt 或参数
(defn with-retry [f {:keys [max-attempts base-ms]}]
  (loop [n 1]
    (let [r (try {:ok (f)} (catch Exception e {:err e}))]
      (cond
        (:ok r) (:ok r)
        (>= n max-attempts) (throw (:err r))
        (retryable? (:err r)) (do (Thread/sleep (* base-ms (bit-shift-left 1 (dec n))))
                                  (recur (inc n)))
        :else (throw (:err r))))))

超时要单独设:LLM 请求可能长达数十秒,必须把连接超时和读取超时分开配置,否则慢请求会占满线程池。

3. 流式处理

3.1 SSE 流解析

LLM 的流式响应是 SSE 格式,每行 data: {...},以 data: [DONE] 结束:

(ns app.llm.stream
  (:require [clojure.string :as str]
            [cheshire.core :as json]
            [clojure.core.async :as a]))

(defn parse-sse-line [line]
  (when (str/starts-with? line "data: ")
    (let [payload (subs line 6)]
      (when-not (= payload "[DONE]")
        (json/parse-string payload true)))))

(defn sse->chan [reader]
  (let [out (a/chan 32)]
    (a/thread
      (try
        (with-open [r reader]
          (doseq [line (line-seq r)]
            (when-let [evt (parse-sse-line line)]
              (a/>!! out evt))))
        (finally (a/close! out))))
    out))

3.2 core.async 编排

把流式增量抽出来,用 transducer 做转换:

(defn text-deltas [ch]
  (a/pipe ch (a/chan 32 (comp
                          (map #(get-in % [:choices 0 :delta :content]))
                          (remove nil?)))))

为什么用 core.async 而不是回调:流式 + 并发是天然需要背压的场景。通道的缓冲区大小就是背压阈值,慢消费者会自然地让上游阻塞,不会无限堆内存。更深入的管道模式见 Clojure core.async 深入 。

3.3 增量渲染与背压

服务端把增量转发给前端(SSE 或 WebSocket)时要注意:

  • 不要每个 token 都推一次:网络开销会压垮连接,应做微批(每 20~50ms 合并一次);
  • 缓冲要有上限:客户端读得慢时必须丢弃或断开,而不是无限缓冲;
  • 首字节时间要监控:TTFT(time to first token)是流式体验的核心指标。
(defn batched [ch ms]
  (let [out (a/chan 8)]
    (a/go-loop [buf []]
      (let [[v port] (a/alts! [ch (a/timeout ms)])]
        (cond
          (= port ch) (if (nil? v)
                        (do (when (seq buf) (a/>! out (apply str buf))) (a/close! out))
                        (recur (conj buf v)))
          :else (do (when (seq buf) (a/>! out (apply str buf))) (recur [])))))
    out))

4. 工具调用与结构化输出

4.1 function calling 协议

模型返回 tool_calls,每个含 name 与 arguments(JSON 字符串)。核心循环是「调用 → 执行 → 回灌结果 → 再调用」:

(defn run-tools [client messages tools max-turns]
  (loop [msgs messages, turn 0]
    (let [resp (p/complete client {:model "gpt-4o" :messages msgs :tools tools})]
      (if-let [calls (:tool_calls resp)]
        (if (< turn max-turns)
          (let [results (mapv execute-tool calls)]
            (recur (into msgs (cons resp results)) (inc turn)))
          (throw (ex-info "tool loop exceeded" {:max max-turns})))
        resp))))

4.2 用 Malli 定义工具 schema

工具的参数契约应该用 schema 表达,既生成给模型的 JSON Schema,又用于本地校验:

(def OrderQuery
  [:map
   [:user-id :string]
   [:status [:enum :pending :paid :shipped]]
   [:limit {:optional true} [:int {:min 1 :max 100}]]])

(def tools
  [{:type "function"
    :function {:name "query_orders"
               :description "按用户与状态查询订单"
               :parameters (json-schema/transform OrderQuery)}}])

双重收益:一份 schema 同时喂给模型(作为参数说明)和本地(作为入参校验)。模型幻觉出来的字段会在执行前被挡掉。

4.3 参数校验与执行

(defn execute-tool [{:keys [function]}]
  (let [{:keys [name arguments]} function
        args (json/parse-string arguments true)]
    (try
      (if (m/validate OrderQuery args)
        {:role "tool" :tool_call_id (:id function)
         :content (json/generate-string (dispatch-tool name args))}
        {:role "tool" :tool_call_id (:id function)
         :content (json/generate-string
                   {:error "invalid arguments"
                    :details (me/humanize (m/explain OrderQuery args))})})
      (catch Exception e
        {:role "tool" :tool_call_id (:id function)
         :content (json/generate-string {:error (.getMessage e)})}))))

关键:工具执行失败要作为结果回灌,而不是抛异常中断。把错误信息告诉模型,它往往能自我修正参数。

4.4 结构化输出的可靠性

  • JSON mode / Structured Outputs:优先用供应商提供的强制 JSON 能力;
  • 本地兜底校验:无论供应商承诺什么,本地必须用 schema 校验一次;
  • 失败重试:把校验错误作为新消息回灌,让模型修正(通常 1~2 次即可)。

5. 向量检索与 RAG

5.1 分块与向量化

(defn chunk [text {:keys [size overlap]}]
  (->> (partition-all (- size overlap) size text)
       (map #(apply str %))))

(defn index! [db docs]
  (let [chunks (mapcat (fn [{:keys [id text meta]}]
                         (map #(hash-map :doc-id id :text % :meta meta)
                              (chunk text {:size 800 :overlap 100})))
                       docs)
        vectors (p/embed client (map :text chunks))]
    (insert-vectors! db (map #(assoc %1 :embedding %2) chunks vectors))))

5.2 检索与重排

(defn retrieve [db query k]
  (let [qv (first (p/embed client [query]))]
    (->> (search-vectors db qv (* k 4))     ; 先粗召回
         (rerank query)                      ; 再精排
         (take k))))

两阶段检索(召回 + 重排)几乎是标配:粗召回用向量相似度保证召回率,重排用更贵的模型(或交叉编码器)提升精度。RAG 的整体架构见 AI LLM RAG 模式 。

5.3 上下文拼装

(defn build-prompt [query hits]
  [{:role :system :content "只依据给定资料回答,资料不足时明确说明。"}
   {:role :user :content
    (str "资料:\n"
         (str/join "\n---\n" (map #(str "[来源 " (:doc-id %) "]\n" (:text %)) hits))
         "\n\n问题:" query)}])

两个工程要点:

  • 来源要带 id:方便引用与溯源,也让模型更容易「有据可依」;
  • 控制总 token:拼装前要按预算裁剪,超出预算就减少召回条数或截断长文档。

6. 成本与并发控制

6.1 token 计量与预算

(defn estimate-tokens [text] (long (/ (count text) 2.5)))  ; 粗略估算

(defn budget-check [usage {:keys [max-tokens-per-call daily-budget]}]
  (when (> (+ (:used usage) (:estimated usage)) daily-budget)
    (throw (ex-info "daily budget exceeded" {:usage usage}))))

生产必备:按调用、按用户、按天三层限额。没有预算控制的 LLM 应用,一次循环 bug 就能烧掉一个月预算。

6.2 并发限流

(defn limited-client [client n]
  (let [sem (java.util.concurrent.Semaphore. n)]
    (reify p/LLMClient
      (complete [_ req]
        (.acquire sem)
        (try (p/complete client req)
             (finally (.release sem))))
      (stream [_ req] (p/stream client req))
      (embed  [_ xs]  (p/embed client xs)))))

用信号量限制并发数,避免触发供应商限流,也避免本地线程池耗尽。更细粒度的限流可以按模型分别配置。

6.3 缓存与去重

(defn cached-complete [client cache req]
  (let [k (hash req)]
    (if-let [hit (get @cache k)]
      hit
      (let [resp (p/complete client req)]
        (swap! cache assoc k resp)
        resp))))

两层缓存:

  1. 精确缓存:请求完全相同直接返回(适合评测与重放);
  2. 语义缓存:query 向量相似度超阈值则复用(适合客服问答类)。

语义缓存能显著降本,但要用相似度阈值 + 人工抽检保证不会答非所问。

7. 可观测与评测

7.1 必埋的指标

指标说明
请求量 / 错误率按模型、按调用方分维度
延迟 P50/P95区分一次性与流式(TTFT)
token 用量输入/输出分开计
成本按天、按功能模块聚合
工具调用成功率参数校验失败率尤其重要
缓存命中率精确与语义分开

7.2 评测与回放

不可变数据让回放变得简单:把每次请求的 messages、model、params 与响应存下来,就是一份可重放的评测集。

(defn replay [client cases]
  (for [{:keys [input expected]} cases]
    (let [actual (p/complete client input)]
      {:input input
       :pass? (= expected (:content actual))
       :actual actual})))

建议:上线前建一个 50~200 条的小评测集,每次改 prompt 或换模型都跑一遍。评测方法论的完整讨论见 AI LLM 应用开发 。

8. 小结

Clojure 集成 LLM 的工程要点,可以压缩成六条:

  1. 先抽象接口再谈实现:供应商无关的 request 数据结构,换模型不改业务;
  2. 流式用 core.async:通道缓冲区就是背压,微批降低网络开销;
  3. 工具参数用 schema 双用:一份 Malli schema 既喂模型又做本地校验;
  4. RAG 走两阶段检索:粗召回保召回率,重排保精度,拼装前控 token;
  5. 成本三层限额:按调用、按用户、按天,缓存分精确与语义两层;
  6. 可观测与评测先行:没有评测集的 prompt 迭代等于盲改。

Clojure 在这里的价值不是「更快的推理」,而是把 LLM 调用当作普通的数据管道来编排——流、并发、重试、缓存、评测,全是它擅长的领域。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「clojure」更多文章

  1. Clojure 桌面 UI:cljfx 与 JavaFX 实战
  2. JSON/EDN 序列化与数据格式互操作
  3. JVM 调优与容器化部署:GC、JFR 与 Docker