博客

从 Kafka 到 Databend Cloud:万亿级 Agent Trace 接入链路的工程实践

avatarJeremy8月 12, 2026
从 Kafka 到 Databend Cloud:万亿级 Agent Trace 接入链路的工程实践

从案例说起:当 Agent Trace 成为 Evals 的数据底座

某中国头部 AI 大模型企业在开源万亿级强化推理大模型后,开始面对一个典型的 Agent 时代数据工程问题。该模型面向百万级超长上下文与 Agent 级复杂任务处理,一次长任务可能包含数千次工具调用,累计处理数百万 context tokens。随着模型能力从单轮对话走向长程任务执行,运行过程中会持续产生请求参数、模型输入输出、Token 用量、首字延迟、工具调用、检索过程、错误信息,以及多层嵌套 Span。

这些 Trace 不只是排查问题的日志。它们记录了模型为什么成功、在哪里失败,以及 Prompt、工具、Harness 和模型版本的变化如何影响最终结果,是构建 Evals 数据集、开展归因分析和推动模型迭代的重要数据资产。

在生产高峰期,该企业的在线 Trace 写入达到每小时 TB 级,并长期累积到万亿级规模。他们选择 Databend Cloud 承载全链路 Agent Trace 数据,构建一条贯穿原始数据接入、增量处理、业务建模、Evals 生产与归因分析的统一数据管道。

Agent Trace 接入难在哪里?

在 Agent Trace 数据管道中,最先承受压力的是接入链路。

一次模型调用会产生大量事件。复杂 Agent 任务还会继续展开为多层 Span,并跨越较长时间窗口。业务高峰到来时,大量 Trace 同时涌入 Kafka,被分散到多个 Topic 和 Partition。接入端需要同时满足几个要求:

  • 尽快追平 Kafka lag,不能让积压持续扩大。

  • 不能因为追求吞吐而提前提交 offset,造成数据丢失。

  • 不能把复杂 JSON 解析全部放在 Consumer 里,否则 Schema 频繁变化会拖慢链路演进。

  • 不能用一条消息触发一次数据库写入,否则固定开销会被放大。

  • 需要在故障、重启、缩容、网络抖动时保持可恢复。

bend-ingest-kafka
工作在这个位置。它从 Kafka 持续消费原始 Trace,以批量文件写入 Databend Cloud 的 Raw Table,再由 Databend Stream 跟踪增量,交给
MERGE INTO
解析、去重并写入最终业务表。

整条链路可以概括成下面这张图:

这套架构的核心边界很清晰:

bend-ingest-kafka
负责把原始数据快速、完整、可恢复地送入 Databend Cloud;复杂的业务解析、字段抽取和去重留在 Databend 内部完成。

这种分工让接入与加工解耦。流量增长时,可以独立扩容 Consumer 和 Databend Warehouse;Trace 字段变化时,也不需要频繁修改 Kafka 消费程序。

把 Kafka 的并行度延伸到 Databend 写入端

Kafka 已经用 Topic 和 Partition 把海量 Trace 切成许多可以并行处理的数据流。要真正利用这份并行能力,消费端不能停留在单进程、单线程模型。

生产环境会部署多个

bend-ingest-kafka
实例。同一 Topic 对应的实例使用相同 consumer group。每个实例内部还可以启动多个 Worker,每个 Worker 都拥有独立的 Kafka Consumer、Batch Reader 和待写入文件队列。Kafka Consumer Group 负责在所有实例和 Worker 之间分配 Partition。

源码中的实例内并发结构很直接。程序按照

workers
配置创建多个 Worker,并让它们各自在 Goroutine 中运行:

wg.Add(cfg.Workers)
for i := 0; i < cfg.Workers; i++ {
w := NewConsumeWorker(cfg, fmt.Sprintf("worker-%d", i), ig)
go func() {
w.Run(ctx)
wg.Done()
}()
}

这不只是多开几个 Kafka Consumer。数据读取之后的 NDJSON 生成、Zstd 压缩、Stage 上传和

COPY INTO
同样会随 Worker 并行展开。实例数和 Worker 数增加后,整个写入通道一起变宽,而不是把压力集中到某个单线程写入点。

有效并行度仍然受 Kafka Partition 数量约束。假设一个 Topic 有 144 个 Partition,部署 12 个实例、每个实例运行 4 个 Worker,就会形成 48 个消费单元,平均每个 Worker 处理约 3 个 Partition。继续增加实例时,Kafka 会通过 rebalance 重新分配 Partition;在线增加 Partition 后,

bend-ingest-kafka
也会定期刷新 Topic metadata,发现变化并重新加入分配。

相关过程会写入日志:

Partitions revoked: [...]
Partitions assigned: [...]

扩容是否生效、新增 Partition 是否已经开始消费,都可以从 rebalance 日志和 Kafka lag 中直接确认。

Raw Mode:先完整接住变化,再做业务建模

这条生产链路使用

bend-ingest-kafka
的 Raw Mode:

{
"isJsonTransform": false
}

在这个模式下,接入程序不会把 Trace JSON 强行映射到固定业务表,而是将原始数据连同 Kafka 元数据一起写入 Raw Table:

CREATE TABLE trace_raw (
uuid STRING,
koffset BIGINT,
kpartition INT,
raw_data JSON,
record_metadata JSON,
add_time TIMESTAMP
);

一条落入 Raw Table 的记录大致如下:

{
"uuid": "0bc7efb9-8e4c-4a8a-a0c9-6a442b0523fe",
"koffset": 18273645,
"kpartition": 37,
"record_metadata": {
"topic": "llm-traces",
"partition": 37,
"offset": 18273645,
"key": "trace-key",
"create_time": "2026-08-10T10:00:00Z"
},
"add_time": "2026-08-10T10:00:01Z",
"raw_data": {
"trace_id": "trace-001",
"span_id": "span-001",
"model": "large-model-v3",
"usage": {
"prompt_tokens": 1234,
"completion_tokens": 356
}
}
}

Agent Trace 的结构变化很快。模型版本更新后,可能新增推理 Token、缓存命中、计费字段;Agent 框架升级后,工具调用、上下文压缩和错误结构也可能发生变化。如果接入程序需要逐字段理解所有数据,上游的一次字段变更就可能变成一次消费者升级。

Raw Mode 把这个问题移出了实时接入的关键路径。只要消息仍是合法 JSON,

bend-ingest-kafka
就先把它完整保存下来。字段如何解析、类型如何转换、最终进入哪张分析表,交给后续 Databend SQL 完成。

这样做有三个收益:

  • 接入端 CPU 消耗更低,不被复杂 JSON 解析拖慢。

  • Schema 演进不会阻塞 Kafka 消费。

  • 原始 Trace 被完整保留,后续解析规则调整后仍然可以重新加工。

Raw Table 还保留了

topic
partition
offset
。这三个字段组合起来,就是一条 Kafka 消息稳定的身份标识。发生故障重放时,下游可以据此去重;需要回溯数据时,也能把 Databend 中的记录准确对应回 Kafka 消费位置。

批量文件写入:不要让一条消息触发一次入仓

海量小请求是分析型数据库写入最不经济的方式。

bend-ingest-kafka
从 Kafka 读取消息后,会先组成 Batch,再生成 NDJSON 文件。

流量充足时,消息数达到

batchSize
,Batch 立即结束;流量较低时,
batchMaxInterval
会限制等待时间,避免数据因为迟迟攒不满而长时间停留。

{
"batchSize": 10000,
"batchMaxInterval": 10
}

高流量 Topic 可以依靠大 Batch 摊薄固定开销,低流量 Topic 也能守住可见延迟。

Raw Trace 里有大量重复的 JSON 字段名,压缩收益通常很可观。开启

copyIntoUploadCompression
后,程序用带缓冲的 Writer 一边生成 NDJSON,一边完成 Zstd 压缩,最后上传的是
.ndjson.zst
文件:

zstdWriter, err := zstd.NewWriter(outputFile)
writer = zstdWriter
bufferedWriter := bufio.NewWriter(writer)

Databend 在

COPY INTO
中使用
COMPRESSION = AUTO
自动识别并解压。网络传输和 Stage 临时存储体积随之下降,特别适合原始 Trace 这类字段重复度较高的 JSON 数据。

Batch 解决了逐条写入的问题,但如果每生成一个文件就执行一次

COPY INTO
,固定开销仍然会被频繁支付。SQL 请求、任务调度、文件解析初始化和事务提交,与文件里有多少条数据并不是完全成正比。文件较小时,这些固定成本尤其明显。

因此,

bend-ingest-kafka
增加了文件数与等待时间两个触发条件:

{
"copyIntoFileCount": 128,
"copyIntoMaxInterval": 5
}

它们是 OR 关系:

  • 文件照常生成后立即上传到 Stage。

  • Pending 文件达到 128 个时立即 COPY,不必等待 5 秒。

  • 从第一个文件上传成功开始达到 5 秒时,即使不足 128 个也立即 COPY。

这样既能在高流量下充分聚合文件,也能限制低流量场景的写入延迟。

最终生成的 SQL 类似这样:

COPY INTO trace_raw
FROM @~/batch/
FILES = (
'file-1.ndjson.zst',
'file-2.ndjson.zst',
'file-3.ndjson.zst',
'file-4.ndjson.zst',
'file-5.ndjson.zst'
)
FILE_FORMAT = (
TYPE = NDJSON
MISSING_FIELD_AS = FIELD_DEFAULT
COMPRESSION = AUTO
)
PURGE = TRUE
FORCE = FALSE
DISABLE_VARIANT_CHECK = TRUE;

假设

batchSize
为 10,000:高流量下,pending 文件达到
copyIntoFileCount
就立即执行 COPY;低流量下,第一个已上传文件等待达到
copyIntoMaxInterval
也会触发 COPY。单文件大小不会失控,固定开销得到摊薄,写入延迟也有明确上限。

每个 Worker 都维护自己的文件队列,不同 Worker 不会把 pending 文件混在一起。多个 Worker、多个实例可以同时生成文件、上传 Stage,并向 Raw Table 发起相互独立的 COPY。Kafka 侧的多 Partition 并发,由此一路延伸到了 Databend Cloud 写入端。

Offset 提交:必须等数据真正进入 Raw Table

接入链路里最容易被忽视,也最容易造成数据丢失的,是 offset 提交时机。

bend-ingest-kafka
的处理顺序始终是:

文件已经上传到 Stage,并不意味着消息已经处理完成。假如这时提前提交 offset,进程却在

COPY INTO
之前退出,Kafka 会认为消息已经消费,而 Raw Table 里还没有数据,最终留下无法自动恢复的缺口。

所以,Worker 在内存中同时保留 Stage 文件和对应的 Kafka Batch。只有整组文件成功写入 Raw Table,才依次提交这些 Batch 的 offset。COPY 失败时,程序继续使用已经上传的同一组文件重试,不会重复生成和上传;这组 pending 文件处理完成之前,Worker 也不会继续读取更多消息。

这意味着下游异常时,积压留在 Kafka,而不是无边界地堆进进程内存。

这条链路采用的是

at-least-once
语义。它优先保证数据不丢,但在一个特殊边界上可能产生重复:Databend 已经完成 COPY,客户端还未来得及提交 Kafka offset,进程就发生故障。重启后,这批消息会再次消费,Raw Table 里可能出现两份相同数据。

这是有意保留的容错空间。后续

MERGE INTO
会按照业务主键聚合同一批 Stream 增量,并结合
action
取出需要落入目标表的最新状态。相比为了追求接入端“恰好一次”而引入复杂的跨系统事务,这种方式更适合高吞吐链路,也更容易在故障后恢复。

网络抖动同样不会立刻打断消费程序。Stage 上传和 COPY 失败后,程序按指数退避重试:

1 秒 → 2 秒 → 4 秒 → 8 秒 → …… → maxRetryDelay

正常发布或实例缩容时,优雅退出会接住最后一个边界。即使 pending 队列里不足

copyIntoFileCount
,Worker 也会在退出前把剩余文件写入 Raw Table、提交对应 offset,然后再关闭 Kafka Consumer。

节点断电、OOM 或

SIGKILL
无法执行这段流程,但未提交的消息仍会由 Kafka 重放,最终再由下游 MERGE 消除重复。

用 Stream 和 MERGE INTO 消除重复写入影响

数据写进 Raw Table 后,接入阶段就结束了。Trace 的字段解析、类型转换和去重不再占用 Kafka Consumer 的资源,而是在 Databend 内部继续完成。这些步骤可以在 Databend Cloud 中拆成多个 Task 执行:前置 Task 消费 Stream 增量并运行

MERGE INTO
,后续 Task 继续完成字段展开、数据清洗和派生指标计算;每个 Task 都可以绑定独立 Warehouse,并按计划或依赖关系自动运行。

如果每次加工都扫描整张 Raw Table,历史数据达到百亿、万亿级后,成本会越来越高。Databend Stream 会跟踪 Raw Table 的增量变化,让下游任务只读取上次处理之后新增的数据。

概念上,可以在 Raw Table 上创建 Append-only Stream:

CREATE STREAM trace_raw_stream
ON TABLE trace_raw
APPEND_ONLY = TRUE;

Stream 后面的

MERGE INTO
会从
raw_data
中投影目标表字段,并在同一批 Stream 增量中按业务主键取最新一条,再根据
action
完成更新、删除或插入。下面的 SQL 与项目中
GenerateMergeIntoSQL
的生成逻辑一致;真实列名和业务主键会随 Trace Schema 调整:

MERGE INTO trace_detail a
USING (
SELECT
raw_data:id AS id,
raw_data:trace_id AS trace_id,
raw_data:span_id AS span_id,
raw_data:model AS model,
raw_data:prompt_tokens AS prompt_tokens,
raw_data:completion_tokens AS completion_tokens,
action
FROM trace_raw_stream
QUALIFY ROW_NUMBER()
OVER (PARTITION BY id ORDER BY add_time) = 1
) b
ON a.id = b.id
WHEN MATCHED AND b.action = 'update' THEN UPDATE *
WHEN MATCHED AND b.action = 'delete' THEN DELETE
WHEN NOT MATCHED AND b.action != 'delete' THEN INSERT *;

同一个业务主键在一批增量里出现多次时,

ROW_NUMBER()
配合
QUALIFY
按项目中的生成逻辑筛出一条记录。
action='update'
更新目标记录,
action='delete'
删除目标记录,其他非删除事件在目标表不存在时插入。接入端的
at-least-once
,由此转化为业务表上的最终幂等。

这也解释了为什么整条链路没有在

bend-ingest-kafka
里解析所有 Trace 字段。接入程序只做稳定且通用的工作,Databend 负责它更擅长的增量计算和结构化处理。Trace 增加字段时,可以修改后面的 SQL,而不必升级全部 Kafka Consumer;解析规则发生错误时,Raw Table 中的原始数据仍然保留,也能按新规则重新加工。

用真实链路验证多文件 COPY 的行为

为了验证 Raw Mode 和多文件 COPY 的实际行为,项目加入了一组 10 万条数据的全链路测试。测试连接真实 Kafka、真实 Databend Stage 和真实 Raw Table,使用下面的参数:

消息总数: 100,000
isJsonTransform: false
batchSize: 1,000
copyIntoFileCount: 5
Worker: 1
Kafka Partition: 1

100,000 条消息按每批 1,000 条生成 100 个文件;每 5 个文件执行一次 COPY,理论上应当产生 20 次

COPY INTO
。实际运行结果与这个关系完全一致:

Raw Table 行数: 100,000
唯一 Kafka offset: 100,000
offset 范围: 0 ~ 99,999
Stage 上传次数: 100
COPY INTO 次数: 20
写入错误: 0

本地测试环境中,单 Worker 完成 Kafka 消费、Raw 数据封装、Zstd 压缩、Stage 上传、COPY 和 offset commit,总耗时约 6.58 秒,吞吐约为 15,208 rows/s。

这组数据用于验证全链路正确性和多文件 COPY 的执行效果,不代表 Databend Cloud 的性能上限。生产吞吐还取决于 Trace 平均大小、Kafka Partition 数、实例与 Worker 数、网络带宽以及 Databend Warehouse 规格。

不过,它清楚地验证了一个关键关系:100,000 条消息生成 100 个文件,最终只需要 20 次 COPY,Raw Table 中仍然得到 100,000 条 offset 连续、元数据完整的记录。

万亿级不是一台机器跑出来的

万亿级描述的是长期累积的数据规模,以及持续面对高峰流量时的系统能力。它不依赖某个配置惊人的单机进程,而是来自整条链路每一层都能扩展。

Kafka 可以拆分更多 Topic、增加 Partition;

bend-ingest-kafka
可以增加实例和 Worker;Databend Cloud 可以根据写入和加工压力调整 Warehouse。Raw Table 与业务表之间由 Stream 解耦后,接入速度也不再和复杂 JSON 解析、去重任务绑定在一起。

流量增长时,这条路径不需要改变系统的基本形态。新的 Consumer 接管更多 Partition,更多 Worker 并行生成和上传文件,多文件 COPY 继续降低固定开销;数据进入 Raw Table 后,再由独立的计算资源完成增量加工。

从 Kafka 看,

bend-ingest-kafka
是一个 Consumer;从 Databend Cloud 看,它是一条持续运行的批量写入通道。真正让它能够承担海量 Trace 接入的,并不是某一个孤立优化,而是这些细节共同形成的结果:

  • Batch 控制单文件大小。

  • Zstd 降低传输和临时存储成本。

  • 多文件

    COPY INTO
    摊薄调度开销。

  • offset 延迟提交守住数据完整性。

  • 重试和优雅退出处理故障边界。

  • 多实例与多 Worker 把 Kafka 的 Partition 并行度延伸到 Databend 写入端。

Raw Table、Stream 和

MERGE INTO
则补上了链路的后半程。原始 Trace 先被完整保存,新增数据随后进入增量处理,重复消息在业务表落地前被消除,最终服务于 Evals、归因分析和大模型可观测查询。

当在线 Trace 写入达到每小时 TB 级,长期数据走向万亿规模,这种“接入先求稳、加工留在后端”的分工,会比一个包办所有逻辑的 Consumer 更容易扩展,也更容易演进。

小结

bend-ingest-kafka
的设计并不追求在 Consumer 里完成所有工作。它把最关键、最通用的接入职责做好:并行消费、批量文件写入、压缩、Stage 上传、多文件 COPY、offset 延迟提交、失败重试和优雅退出。

Databend Cloud 则接过后半程:Raw Table 保留事实,Stream 追踪增量,

MERGE INTO
完成结构化处理与最终幂等,Warehouse 隔离写入、加工和分析负载。

对 Agent Trace 来说,数据平台的价值不只是“把日志存下来”。更关键的是,它要把持续增长、结构漂移、包含敏感信息的原始执行记录,转化成可以被工程团队反复使用的数据资产。Databend Cloud 在这条链路中承担的是 Agent Trace 的统一数据层:Raw Table 保留完整 JSON 与 Kafka 元数据,Stream 和

MERGE INTO
支撑面向 Evals 的增量加工,独立 Warehouse 隔离写入、加工和分析负载,完整 Trace 则可以继续服务归因、回放、训练数据生产和模型迭代。

Databend Cloud 不只是分析型数据库,而是面向现代 AI 应用的 Agent Trace 数据底座。它让 Trace 从“事后排查日志”变成“持续反馈系统”,帮助团队更快定位问题、更快构建评测数据集,也更快推动模型和 Agent 产品进化。

对于正在建设 Agent Trace 数据管道的团队,这套架构提供了一个可复用的工程思路:不要让接入层理解所有业务语义,先完整、稳定、可恢复地接住数据,再把复杂加工交给更适合做增量计算和分析的数据平台。

bend-ingest-kafka
项目已开源:GitHub:https://github.com/databendcloud/bend-ingest-kafka

相关阅读:

《从万亿级大模型到全线应用:Databend Cloud 助力头部 AI 企业构建全链路 Trace 数据管道》

《Agent 轨迹分析与归因的数据工程实践》

分享本篇文章

订阅我们的新闻简报

及时了解功能发布、产品规划、支持服务和云服务的最新信息!