把大模型 API 直接写进同步 Web 请求,原型阶段很简单:收到问题、调用模型、返回答案。但进入真实业务后,模型响应时间并不完全可控,还可能遇到连接超时、上游限流、临时服务错误或调用方主动取消。若 Web 进程一直等待,连接池、工作线程和反向代理超时会相互影响;若调用失败后立即重试,又可能形成流量放大。
更稳妥的做法是把“接受任务”和“执行模型调用”拆开:业务接口只校验请求、生成任务标识并写入 Kafka;工作进程按自身容量消费消息;结果写入数据库或对象存储,再由客户端轮询或通过回调获取。Kafka 在这里不是模型网关,而是任务缓冲区和可回放的交付日志。
这种架构适合允许异步返回的摘要、分类、文档分析和批量生成任务。对必须在数百毫秒内完成的交互,它未必合适,因为排队、序列化和结果查询都会增加额外延迟。
核心原理
一条最小链路包含四个角色:
- API 服务生成全局唯一的
task_id,先登记任务,再向llm.jobs发布消息。 - Kafka 按消息键分区。使用
task_id作为键,可以让同一任务的相关消息稳定进入同一分区,但不能自动实现业务幂等。 - Worker 拉取消息,调用模型接口,将成功结果持久化后再提交消费位点。
- 超过重试上限的消息写入
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.30s、llm.jobs.retry.5m,由专门调度器按到期时间重新投递;Kafka 本身不是通用定时任务系统,延迟语义需要应用层实现。
可观测与安全边界
至少记录 task_id、主题、分区、位点、消费组、尝试次数、模型标识、耗时和归一化错误类型。提示词与模型原文默认不应进入普通日志;确需留存时,应设置访问控制、脱敏和保留期限。
建议同时监控:消费积压、任务等待时间、成功率、重试率、死信数量和单任务处理时间。积压上升不一定意味着 Kafka 故障,也可能是上游变慢、Worker 容量不足或限流策略生效。扩容前要检查模型侧并发与配额,避免消费者增加后把压力直接转移给上游。
常见问题
为什么关闭自动提交?
自动提交可能在业务结果尚未持久化时推进位点。手动提交让“何时确认消息”与业务完成状态一致,但仍须处理重复消费。
能否保证每个任务只调用一次模型?
通常不能绝对保证。Kafka 消费和本地存储可以通过幂等降低重复结果,但外部 HTTP 调用不在同一事务中。若供应商没有可靠的幂等键机制,就应按“可能重复调用”设计预算、审计与补偿。
为什么不把失败消息一直留在原主题?
持续失败的消息会占用处理能力,并让排障缺少明确入口。死信主题可以隔离无法自动恢复的任务,但必须配置告警和处置流程,否则它只是另一个无人查看的积压队列。
Kafka 可以直接暴露到公网吗?
不建议为了远程调试直接开放 broker 端口。应优先使用私网、VPN 或受控隧道,并配置 TLS、认证和 ACL。advertised.listeners 必须发布客户端实际可达的地址;具体写法依部署方式和安全协议而异。
如何避免不同模型互相拖累?
可按风险与容量拆分主题或消费组,例如交互任务与批处理任务分离。不要只按供应商名称拆分,否则业务优先级、数据等级和重试策略仍可能混在一起。
总结
Kafka 能把模型调用从同步请求中解耦出来,但可靠性并非来自“加一个消息队列”,而是来自完整的失败语义:任务标识、手动提交、结果幂等、有限重试、死信处置、双写补偿和可观测指标。先用最小链路验证消息契约和故障恢复,再根据实际积压、上游限制与合规要求增加延迟主题、Outbox 和多级消费者,才能让异步模型任务既可扩展,也可追踪和恢复。