从案例说起:当 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
MERGE INTO
整条链路可以概括成下面这张图:

这套架构的核心边界很清晰:
bend-ingest-kafka
这种分工让接入与加工解耦。流量增长时,可以独立扩容 Consumer 和 Databend Warehouse;Trace 字段变化时,也不需要频繁修改 Kafka 消费程序。
把 Kafka 的并行度延伸到 Databend 写入端
Kafka 已经用 Topic 和 Partition 把海量 Trace 切成许多可以并行处理的数据流。要真正利用这份并行能力,消费端不能停留在单进程、单线程模型。
生产环境会部署多个
bend-ingest-kafka

源码中的实例内并发结构很直接。程序按照
workers
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
有效并行度仍然受 Kafka Partition 数量约束。假设一个 Topic 有 144 个 Partition,部署 12 个实例、每个实例运行 4 个 Worker,就会形成 48 个消费单元,平均每个 Worker 处理约 3 个 Partition。继续增加实例时,Kafka 会通过 rebalance 重新分配 Partition;在线增加 Partition 后,
bend-ingest-kafka
相关过程会写入日志:
Partitions revoked: [...]
Partitions assigned: [...]
扩容是否生效、新增 Partition 是否已经开始消费,都可以从 rebalance 日志和 Kafka lag 中直接确认。
Raw Mode:先完整接住变化,再做业务建模
这条生产链路使用
bend-ingest-kafka
{
"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
这样做有三个收益:
-
接入端 CPU 消耗更低,不被复杂 JSON 解析拖慢。
-
Schema 演进不会阻塞 Kafka 消费。
-
原始 Trace 被完整保留,后续解析规则调整后仍然可以重新加工。
Raw Table 还保留了
topic
partition
offset
批量文件写入:不要让一条消息触发一次入仓
海量小请求是分析型数据库写入最不经济的方式。
bend-ingest-kafka
流量充足时,消息数达到
batchSize
batchMaxInterval
{
"batchSize": 10000,
"batchMaxInterval": 10
}
高流量 Topic 可以依靠大 Batch 摊薄固定开销,低流量 Topic 也能守住可见延迟。
Raw Trace 里有大量重复的 JSON 字段名,压缩收益通常很可观。开启
copyIntoUploadCompression
.ndjson.zst
zstdWriter, err := zstd.NewWriter(outputFile)
writer = zstdWriter
bufferedWriter := bufio.NewWriter(writer)
Databend 在
COPY INTO
COMPRESSION = AUTO
Batch 解决了逐条写入的问题,但如果每生成一个文件就执行一次
COPY INTO
因此,
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
copyIntoFileCount
copyIntoMaxInterval
每个 Worker 都维护自己的文件队列,不同 Worker 不会把 pending 文件混在一起。多个 Worker、多个实例可以同时生成文件、上传 Stage,并向 Raw Table 发起相互独立的 COPY。Kafka 侧的多 Partition 并发,由此一路延伸到了 Databend Cloud 写入端。
Offset 提交:必须等数据真正进入 Raw Table
接入链路里最容易被忽视,也最容易造成数据丢失的,是 offset 提交时机。
bend-ingest-kafka

文件已经上传到 Stage,并不意味着消息已经处理完成。假如这时提前提交 offset,进程却在
COPY INTO
所以,Worker 在内存中同时保留 Stage 文件和对应的 Kafka Batch。只有整组文件成功写入 Raw Table,才依次提交这些 Batch 的 offset。COPY 失败时,程序继续使用已经上传的同一组文件重试,不会重复生成和上传;这组 pending 文件处理完成之前,Worker 也不会继续读取更多消息。
这意味着下游异常时,积压留在 Kafka,而不是无边界地堆进进程内存。
这条链路采用的是
at-least-once
这是有意保留的容错空间。后续
MERGE INTO
action
网络抖动同样不会立刻打断消费程序。Stage 上传和 COPY 失败后,程序按指数退避重试:
1 秒 → 2 秒 → 4 秒 → 8 秒 → …… → maxRetryDelay
正常发布或实例缩容时,优雅退出会接住最后一个边界。即使 pending 队列里不足
copyIntoFileCount

节点断电、OOM 或
SIGKILL
用 Stream 和 MERGE INTO 消除重复写入影响
数据写进 Raw Table 后,接入阶段就结束了。Trace 的字段解析、类型转换和去重不再占用 Kafka Consumer 的资源,而是在 Databend 内部继续完成。这些步骤可以在 Databend Cloud 中拆成多个 Task 执行:前置 Task 消费 Stream 增量并运行
MERGE INTO
如果每次加工都扫描整张 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
action
GenerateMergeIntoSQL
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

用真实链路验证多文件 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

流量增长时,这条路径不需要改变系统的基本形态。新的 Consumer 接管更多 Partition,更多 Worker 并行生成和上传文件,多文件 COPY 继续降低固定开销;数据进入 Raw Table 后,再由独立的计算资源完成增量加工。
从 Kafka 看,
bend-ingest-kafka
-
Batch 控制单文件大小。
-
Zstd 降低传输和临时存储成本。
-
多文件
摊薄调度开销。COPY INTO -
offset 延迟提交守住数据完整性。
-
重试和优雅退出处理故障边界。
-
多实例与多 Worker 把 Kafka 的 Partition 并行度延伸到 Databend 写入端。
Raw Table、Stream 和
MERGE INTO
当在线 Trace 写入达到每小时 TB 级,长期数据走向万亿规模,这种“接入先求稳、加工留在后端”的分工,会比一个包办所有逻辑的 Consumer 更容易扩展,也更容易演进。
小结
bend-ingest-kafka
Databend Cloud 则接过后半程:Raw Table 保留事实,Stream 追踪增量,
MERGE INTO
对 Agent Trace 来说,数据平台的价值不只是“把日志存下来”。更关键的是,它要把持续增长、结构漂移、包含敏感信息的原始执行记录,转化成可以被工程团队反复使用的数据资产。Databend Cloud 在这条链路中承担的是 Agent Trace 的统一数据层:Raw Table 保留完整 JSON 与 Kafka 元数据,Stream 和
MERGE INTO
Databend Cloud 不只是分析型数据库,而是面向现代 AI 应用的 Agent Trace 数据底座。它让 Trace 从“事后排查日志”变成“持续反馈系统”,帮助团队更快定位问题、更快构建评测数据集,也更快推动模型和 Agent 产品进化。
对于正在建设 Agent Trace 数据管道的团队,这套架构提供了一个可复用的工程思路:不要让接入层理解所有业务语义,先完整、稳定、可恢复地接住数据,再把复杂加工交给更适合做增量计算和分析的数据平台。
bend-ingest-kafka
相关阅读:
《从万亿级大模型到全线应用:Databend Cloud 助力头部 AI 企业构建全链路 Trace 数据管道》
订阅我们的新闻简报
及时了解功能发布、产品规划、支持服务和云服务的最新信息!






