用 Kafka 解耦大模型调用:构建可重试、可追踪的异步任务队列

简介: 本文介绍一种基于Kafka的异步大模型调用架构:将请求接入与模型执行解耦,通过任务ID幂等、手动位点提交、死信队列和结果持久化,保障高可靠性与可观测性,适用于摘要、分类等非实时场景。(239字)

把大模型 API 直接写进同步 Web 请求,原型阶段很简单:收到问题、调用模型、返回答案。但进入真实业务后,模型响应时间并不完全可控,还可能遇到连接超时、上游限流、临时服务错误或调用方主动取消。若 Web 进程一直等待,连接池、工作线程和反向代理超时会相互影响;若调用失败后立即重试,又可能形成流量放大。

更稳妥的做法是把“接受任务”和“执行模型调用”拆开:业务接口只校验请求、生成任务标识并写入 Kafka;工作进程按自身容量消费消息;结果写入数据库或对象存储,再由客户端轮询或通过回调获取。Kafka 在这里不是模型网关,而是任务缓冲区和可回放的交付日志。

这种架构适合允许异步返回的摘要、分类、文档分析和批量生成任务。对必须在数百毫秒内完成的交互,它未必合适,因为排队、序列化和结果查询都会增加额外延迟。

核心原理

一条最小链路包含四个角色:

  1. API 服务生成全局唯一的 task_id,先登记任务,再向 llm.jobs 发布消息。
  2. Kafka 按消息键分区。使用 task_id 作为键,可以让同一任务的相关消息稳定进入同一分区,但不能自动实现业务幂等。
  3. Worker 拉取消息,调用模型接口,将成功结果持久化后再提交消费位点。
  4. 超过重试上限的消息写入 llm.jobs.dlq,由人工或补偿程序检查。

关键点是提交顺序。若先提交位点、后保存结果,进程在两者之间崩溃会造成任务丢失;若先保存结果、后提交位点,崩溃后消息可能再次消费。因此,常见选择是“至少一次消费 + 业务幂等”,接受重复投递,并依靠 task_id 唯一约束阻止重复写结果。

这仍不能绝对消除重复的外部模型调用:Worker 可能已经获得上游响应,却在本地落库前退出。只有当上游支持且明确承诺幂等键语义时,才能进一步压缩这一窗口;否则应把“可能产生重复调用及费用”纳入设计和告警。

准备 Kafka 主题

以下命令假设已经有可访问的 Kafka 集群,kafka-topics.sh 位于 PATH。集群采用 KRaft 还是 ZooKeeper 不影响本文的生产者、消费者设计,但服务端部署参数应以所用发行版文档为准。

export KAFKA_BOOTSTRAP_SERVERS='127.0.0.1:9092'

kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP_SERVERS" \
  --create --if-not-exists --topic llm.jobs \
  --partitions 3 --replication-factor 1

kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP_SERVERS" \
  --create --if-not-exists --topic llm.jobs.dlq \
  --partitions 3 --replication-factor 1

replication-factor=1仅适用于单节点开发环境。生产值必须结合 broker 数量、故障容忍目标和存储成本确定,不能直接照搬。分区数决定消费并行度上限之一,但增加分区也会提高连接、文件和再均衡成本。

安装示例依赖:

python -m pip install confluent-kafka requests

发布任务

生产者只发送模型无关的业务请求,不把 API 密钥放进消息。消息中也应避免携带不必要的个人信息或完整机密文档。

# producer.py
import json
import os
import uuid
from confluent_kafka import Producer

bootstrap = os.environ["KAFKA_BOOTSTRAP_SERVERS"]
producer = Producer({
   "bootstrap.servers": bootstrap})

task_id = str(uuid.uuid4())
job = {
   
    "task_id": task_id,
    "prompt": "用三点概括这份变更说明",
    "attempt": 0
}

producer.produce(
    "llm.jobs",
    key=task_id.encode(),
    value=json.dumps(job, ensure_ascii=False).encode()
)
producer.flush()
print(task_id)

严格场景下,数据库登记任务与 Kafka 发布之间还存在“双写”问题。可采用事务性 Outbox:API 在同一个数据库事务里写入任务表和待发布事件表,再由独立发布器把事件送入 Kafka。Kafka 事务只能协调 Kafka 内部操作,不能自动让普通数据库写入与消息发布成为一个原子事务。

接入模型 API 的 Worker

示例把接口地址、模型名和密钥全部放入环境变量。若服务商当前文档确认其接口与示例请求格式兼容,可设置对应的基础地址;评估中也可以查阅 HaerAPI 的当前文档,但不要仅凭兼容性描述推断数据保留、地域、配额或可用模型。

export MODEL_BASE_URL='https://api.example.com/v1'
export MODEL_API_KEY='从密钥管理系统注入'
export MODEL_NAME='按供应商当前文档填写'
export KAFKA_BOOTSTRAP_SERVERS='127.0.0.1:9092'
python worker.py

下面的示例假设所选接口支持 POST /chat/completions 及相应 JSON 结构。若实际供应商协议不同,应只替换 call_model,不要让供应商字段扩散到 Kafka 消息契约中。

# worker.py
import json
import os
import sqlite3
import time
import requests
from confluent_kafka import Consumer, Producer

MAX_ATTEMPTS = 3
bootstrap = os.environ["KAFKA_BOOTSTRAP_SERVERS"]
consumer = Consumer({
   
    "bootstrap.servers": bootstrap,
    "group.id": "llm-workers-v1",
    "enable.auto.commit": False,
    "auto.offset.reset": "earliest"
})
producer = Producer({
   "bootstrap.servers": bootstrap})
consumer.subscribe(["llm.jobs"])

db = sqlite3.connect("tasks.db")
db.execute("""
CREATE TABLE IF NOT EXISTS task_results (
  task_id TEXT PRIMARY KEY,
  status TEXT NOT NULL,
  result TEXT,
  updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
)
""")


def call_model(prompt):
    response = requests.post(
        os.environ["MODEL_BASE_URL"].rstrip("/") + "/chat/completions",
        headers={
   
            "Authorization": "Bearer " + os.environ["MODEL_API_KEY"],
            "Content-Type": "application/json"
        },
        json={
   
            "model": os.environ["MODEL_NAME"],
            "messages": [{
   "role": "user", "content": prompt}]
        },
        timeout=(5, 60)
    )
    response.raise_for_status()
    body = response.json()
    return body["choices"][0]["message"]["content"]


while True:
    message = consumer.poll(1.0)
    if message is None:
        continue
    if message.error():
        print("kafka_error", message.error())
        continue

    job = json.loads(message.value())
    task_id = job["task_id"]
    exists = db.execute(
        "SELECT 1 FROM task_results WHERE task_id=? AND status='done'",
        (task_id,)
    ).fetchone()
    if exists:
        consumer.commit(message=message, asynchronous=False)
        continue

    try:
        result = call_model(job["prompt"])
        db.execute(
            "INSERT OR REPLACE INTO task_results(task_id,status,result) VALUES(?,?,?)",
            (task_id, "done", result)
        )
        db.commit()
        consumer.commit(message=message, asynchronous=False)
    except (requests.Timeout, requests.ConnectionError) as exc:
        job["attempt"] += 1
        if job["attempt"] >= MAX_ATTEMPTS:
            job["error_type"] = type(exc).__name__
            producer.produce(
                "llm.jobs.dlq", key=task_id.encode(),
                value=json.dumps(job, ensure_ascii=False).encode()
            )
            producer.flush()
            consumer.commit(message=message, asynchronous=False)
        else:
            time.sleep(2 ** job["attempt"])
            producer.produce(
                "llm.jobs", key=task_id.encode(),
                value=json.dumps(job, ensure_ascii=False).encode()
            )
            producer.flush()
            consumer.commit(message=message, asynchronous=False)

这是便于理解的最小实现。它只把连接错误和超时视为可重试错误;HTTP 429、部分 5xx 是否重试,应依据供应商文档和 Retry-After 等响应信息决定。鉴权失败、请求格式错误通常不应盲目重试。

示例中的 sleep 会阻塞当前消费者,也可能在等待过久时影响消费者组稳定性。生产系统更适合建立多个延迟主题,例如 llm.jobs.retry.30sllm.jobs.retry.5m,由专门调度器按到期时间重新投递;Kafka 本身不是通用定时任务系统,延迟语义需要应用层实现。

可观测与安全边界

至少记录 task_id、主题、分区、位点、消费组、尝试次数、模型标识、耗时和归一化错误类型。提示词与模型原文默认不应进入普通日志;确需留存时,应设置访问控制、脱敏和保留期限。

建议同时监控:消费积压、任务等待时间、成功率、重试率、死信数量和单任务处理时间。积压上升不一定意味着 Kafka 故障,也可能是上游变慢、Worker 容量不足或限流策略生效。扩容前要检查模型侧并发与配额,避免消费者增加后把压力直接转移给上游。

常见问题

为什么关闭自动提交?

自动提交可能在业务结果尚未持久化时推进位点。手动提交让“何时确认消息”与业务完成状态一致,但仍须处理重复消费。

能否保证每个任务只调用一次模型?

通常不能绝对保证。Kafka 消费和本地存储可以通过幂等降低重复结果,但外部 HTTP 调用不在同一事务中。若供应商没有可靠的幂等键机制,就应按“可能重复调用”设计预算、审计与补偿。

为什么不把失败消息一直留在原主题?

持续失败的消息会占用处理能力,并让排障缺少明确入口。死信主题可以隔离无法自动恢复的任务,但必须配置告警和处置流程,否则它只是另一个无人查看的积压队列。

Kafka 可以直接暴露到公网吗?

不建议为了远程调试直接开放 broker 端口。应优先使用私网、VPN 或受控隧道,并配置 TLS、认证和 ACL。advertised.listeners 必须发布客户端实际可达的地址;具体写法依部署方式和安全协议而异。

如何避免不同模型互相拖累?

可按风险与容量拆分主题或消费组,例如交互任务与批处理任务分离。不要只按供应商名称拆分,否则业务优先级、数据等级和重试策略仍可能混在一起。

总结

Kafka 能把模型调用从同步请求中解耦出来,但可靠性并非来自“加一个消息队列”,而是来自完整的失败语义:任务标识、手动提交、结果幂等、有限重试、死信处置、双写补偿和可观测指标。先用最小链路验证消息契约和故障恢复,再根据实际积压、上游限制与合规要求增加延迟主题、Outbox 和多级消费者,才能让异步模型任务既可扩展,也可追踪和恢复。

相关文章
|
26天前
|
SQL BI 测试技术
SQL Server 迁移到 KingbaseES:复杂 BI 查询的兼容性验证与性能治理
本文聚焦SQL Server向KingbaseES迁移中复杂BI查询的落地挑战,强调不能仅凭“能执行”判断成功,须同步验证结果语义、执行计划与边界行为。涵盖语法/类型/优化器差异、基线建立、安全改写、自动化校验及执行计划分析,并指出模型可辅助但不可替代数据库实证验证。(239字)
64 0
|
26天前
|
JSON API 数据安全/隐私保护
免费外汇汇率查询接口推荐:官方稳定方案与开源可用清单
本文实测推荐4个免费外汇汇率接口:Frankfurter(ECB数据,免Key、支持1999年起历史)、fawazahmed0(200+币种含加密货币、无速率限制)、open.er-api(160+币种、一行URL获取)、万维易源(官方自营,含K线/转换等多接入点,需appKey)。均经真实连通验证,适配跨境电商、旅行记账与金融学习场景。
398 1
免费外汇汇率查询接口推荐:官方稳定方案与开源可用清单
|
5月前
|
人工智能 缓存 数据中心
大模型应用:大模型多线程推理:并发请求的处理与资源隔离实践.77
本文详解大模型多线程推理与资源隔离技术:通过共享模型、隔离缓存、限制线程数/生成长度/超时时间,实现高并发、低延迟、稳服务。单线程串行耗时85.7秒,多线程(3线程)降至66.5秒,显著提升吞吐量与资源利用率,是大模型规模化落地的核心工程实践。
829 7
|
7月前
|
人工智能 数据可视化 Java
AI智能体的开发方法
本文系统梳理国内AI智能体开发全景:从“感知-决策-行动-记忆”认知闭环架构出发,对比Dify、Coze等低代码平台与LangGraph、AgentScope、Eino、Spring AI Alibaba等编程级框架;解析MCP协议、RAG技术栈等基础设施;并按MVP、企业级、极客定制三类场景给出选型建议。(239字)
|
22天前
|
人工智能 Linux API
Codex接入DeepSeek‑V4‑Flash完整实操:两套方案补齐识图能力保姆级教程
在AI编程Agent工具生态当中,Codex凭借强大的本地工程读写、代码修改、命令执行能力,成为开发者日常项目调试、代码重构、问题排查的常用客户端。DeepSeek‑V4‑Flash作为一款高性价比的文本大模型,拥有超大上下文窗口,Agent任务规划、代码生成、逻辑推演表现十分突出,API调用成本低廉,非常适合作为Codex底层推理基座。但是该模型属于纯文本推理模型,原生并不支持图像输入,当开发者把报错截图、UI界面截图、架构图、数据图表粘贴进会话,模型会直接提示无法解析图片内容,很多开发场景就此被卡住。
336 1
|
22天前
|
人工智能 安全 API
计算巢 X DeepSeek Harness — 云端智能体工作台
DeepSeek Harness是DeepSeek推出的开源云端AI协作者,支持一键部署、多模型切换与安全审批。内置133个插件,可读写代码、执行命令、任务分解、子代理委派,真正“能干活”。团队共享、浏览器即用,告别本地环境噩梦。
计算巢 X DeepSeek Harness — 云端智能体工作台
|
23天前
|
Java Linux 开发工具
IDEA官网下载2026|IDEA社区版安装+使用图文步骤
IntelliJ IDEA(简称IDEA)是JetBrains推出的主流Java集成开发环境,集代码编写、编译、调试、版本控制等功能于一体。提供免费社区版(支持Java/Kotlin/Maven/Gradle等)和付费旗舰版(增强Spring、数据库、前端等支持),开箱即用、智能提示强,适合学习与日常开发。
|
24天前
|
缓存 安全 程序员
智谱GLM-5.3发布同基座纯靠后训练编程涨50还点亮网安技能树
智谱 8 月 14 日发布 GLM-5.3,与 5.2 同基座、纯后训练,编程内部基准提升 50%,CyberGym 拿下开源第一,两周后开源权重
智谱GLM-5.3发布同基座纯靠后训练编程涨50还点亮网安技能树
|
25天前
|
SQL 人工智能 NoSQL
DBX:仅 20MB 的全能数据库工作台,让 70+ 数据库管理变得简单高效
DBX是一款开源轻量级多数据库管理工具,单文件仅20MB,无需Java/Python/Chromium依赖,支持MySQL、达梦、Redis、ClickHouse等70+数据库。集成SQL编辑、数据浏览、AI助手、MCP协议及Docker/Web部署,兼顾高效、安全与智能。
249 0
|
25天前
|
人工智能 自然语言处理 数据可视化
保姆级教程:通义千问最新版全功能详解,零基础上手 + API 调用实操
随着通用人工智能技术持续落地,很多普通用户、办公从业者、学生以及开发者都开始接触大模型工具,但不少新手面对繁杂的模型能力、设置选项、接口调用方式时容易无从下手。通义千问作为通用大模型产品,持续迭代更新,不断补齐长文本处理、多模态解析、自主智能体、工程代码生成、批量文档处理等能力,兼顾网页可视化简易操作与面向开发者的API程序化接入,既适合零基础用户直接在网页端完成办公、学习、内容创作任务,也能满足技术人员搭建AI应用、自动化脚本、知识库系统的开发需求。本文为保姆级完整教程,从基础定位、逐项拆解最新核心功能、分场景实操演示、API代码调用、新手高频避坑要点等维度完整讲解,无论你是完全没有AI使用
603 0