重构一个 50 亿行的事件存储:一路做到引擎里
英文版:English
一段话版本
我们的反欺诈事件层要同时服务两种日益冲突的负载:serving 路径上的毫秒级点查,和越来越多由面向 AI agent 的产品功能(而不是人类分析师)产生的不可预测 ad-hoc 分析。旧的 ClickHouse shared-nothing 层对第二种负载既不能隔离也不能弹性扩容:一条重查询能饿死共享池,我们实测过一条 SELECT * LIMIT 10 排在 bulk insert 后面,墙钟等了 61 秒、CPU 只用了 60 毫秒。我做的不是一次迁移。我把这一层重建在存算分离上(Apache Doris shared-data 模式 + Iceberg-on-S3,运行在 Kubernetes 上),并把查询路由做进了 Doris 前端本身:系统在执行前 6–8 毫秒内给每条查询判轻重,驱动一个 0→N、用完归零的 spot 计算池。迁移对约 52 亿行 / 4 TiB / 最宽约 3,700 列做了 99.945% 行级对账认证;点查延迟从约 8 秒降到约 20 毫秒,并证明在数据增长 19 倍后保持平坦;过程中我根因了一串引擎级故障,从让所有后端在一句 COUNT(*) 上集体 OOM 的 Parquet footer 爆炸,到由 46 倍压缩衰减驱动的 compaction 死亡螺旋。
为什么离开 ClickHouse(它当时还更快)
在我们决定离开时,ClickHouse 在我们自己的一手 benchmark 上仍以 1.2–2.5 倍领先平表扫描。迁移的理由是结构性的,而保持这一点需要纪律,因为「新引擎更快」是所有人都想讲、也最经不起盘问的故事。
四个结构性折中驱动了这个决定:
-
**没有负载隔离。**一条重查询能饿死共享节点:上面那个 61 秒 / 60 毫秒的饥饿案例是日常,不是特例。
-
**刚性的常驻计算。**shared-nothing 把计算绑在本地盘上,空闲期和重查询高峰付同一份常驻硬件的钱。ClickHouse 生态里真正无状态的弹性计算只存在于闭源托管服务中;这是架构缺口,不是配置差异。
-
**脆弱的宽 schema 演进。**每租户 schema 动态演化,落地成数千个稀疏列,靠脆弱的
ALTER ADD COLUMN维护。 -
**用采样交正确性税。**分析查询默认只扫 100 万行采样,因为全量太贵:用正确答案换延迟。
承重论点,也是对最新版 ClickHouse 依然成立的一个,是硬内存隔离:Doris 的 workload group 对每个查询池施加 cgroup 级硬限,而 ClickHouse 的对应机制仍是 best-effort。当一个面向 AI agent 的产品层开始以不可预测的频率生成不可预测的查询形状时,隔离与弹性就再也不是可选项。
目标架构
Doris 4.0.5 shared-data 模式:tablet 数据住在 S3 storage vault,元数据在 FoundationDB 背书的 MetaService,后端是无状态计算,本地 EBS 只当 file cache。计算按 compute group 物理分池:常驻 serving 池在 on-demand 节点,弹性重查询池在 spot 上、静息 0 副本。每个池内再用 workload group 治理 CPU、内存与并发。同一个 S3 上的 Iceberg 充当开放湖层。
值得明说的设计洞察:存算分离没有消灭状态,它把状态收拢进一层(MetaService、FoundationDB、S3 recycler),而那一层恰恰是你最不能偷工减料的地方。三个新的有状态故障域取代了「每个节点挂着盘」,下面的运维工作反映的正是这笔交换。
迁移是一个内存工程问题
约 3,700 列的表会在四个地方打断朴素的「导出 Parquet、批量灌入」路径,没有一个是显然的,而且加硬件一个都治不了:
-
**Footer 炸弹。**Parquet footer 随
row_groups × columns增长。早期对象带着 2.12 GB 的 footer;reader 在碰到任何一行数据之前就 OOM 了全部八台后端。标志性证据:连一句COUNT(*)都死了,证明爆炸的是元数据解析而不是数据扫描。重构对象后 footer 降约 9 倍。 -
**不被追踪的 scanner。**Doris 向量化 Parquet column reader 的内存分配在
exec_mem_limit之外,规模是 scanner 并发 × 列数;默认配置下约每台后端 78 GB。任何查询级内存限制都管不住一块引擎不追踪的缓冲。 -
**导出缓冲的跷跷板。**每条导出流缓冲
row_group_size × columns。切小片救不了导出内存;调小 row group 又会引爆 footer。只有同时切小对象和 row group,跷跷板才解耦。 -
**导入内存随表增长,而不是随批次。**灌到约 17 亿行时,完全相同的 load 任务开始冲破内存上限并 livelock:merge-on-write 的 delete-bitmap 维护和宽表 compaction 随累计表大小增长。一个 watchdog(软停/硬杀阈值、原子中止、幂等重试)稳住了局面;降到单流并发后零数据丢失地解决,吞吐本来就受导出节奏限制。
最终成型的流水线按字节预算调度、逐对象 checkpoint、任何时刻可安全 kill,在四台后端上跑出约 19,500 rows/s:在 3,700 列 schema 上约为小批方案的 50 倍。还有一个我没在任何地方见人发表过的结果:在这种宽度的 schema 上,列宽单独就能让写吞吐摆动约 5 倍,而完整的正确性栈(merge-on-write 去重、三个倒排索引、ZSTD)代价不到 4–12%。看起来昂贵的特性几乎免费,宽度本身才是成本。
负载不是我们以为的样子
原本的 serving 表方案建立在一套手挑的 50 条查询 benchmark 上。在承诺一个不可逆的布局之前,我挖了 14 天的生产 system.query_log(7,470 万行日志),图景整个翻转:流量的绝大多数(远超九成)是点查,按一个高基数用户标识键入,其中大多是无界的「最新值」重建,补不上任何时间过滤。第二张 serving 表从「也许」变成硬需求。照着 benchmark 设计,优化的会是错误的负载。
然后是物理课。Iceberg 布局上的点查跑 7–13 秒,ClickHouse 是 21 毫秒。反射式的修法(bloom filter)砍掉了 73% 的扫描行数,墙钟一动不动:一个活跃用户的行散落在 2,044 个分区文件里,地板是打开文件的数量,不是扫描的行数。(我们还字节级证明了两条常见 Iceberg writer 路径会静默忽略声明的 bloom-filter 表属性。)**索引修不了布局问题。**修法是物理重聚簇:一张按用户键 hash 分布的 Doris 原生表。分布键是唯一日后不能 ALTER 的决策,所以先做了双表 A/B(同一份数据、镜像的两个键、四组测量)再定案。结果:暖读 16–24 毫秒,并在迁移的三个检查点(1,500 万、1.07 亿、2.86 亿行)重测,证明延迟在表增长 19 倍后保持平坦。
路由是准入控制,不是性能调优
serving 和分析同平台后,那不到 1% 的重查询必须在执行前被拦住:扩容是分钟级机制,OOM 是秒级事件,弹性本身永远当不了安全保证。
原生 Doris EXPLAIN 在我们生产形状的 SQL 套件上只有 45% 路由准确率,重查询漏判率 27%。决定性的观察:同一模板换参数后 EXPLAIN 文本逐字节相同,真实内存却差 23 倍。这是信号问题,不是规则问题:优化器内部本来就算好了每个算子的基数、选择率和状态大小,只是从不打印。
于是我在一个 Apache Doris fork 上把它修进引擎:
-
EXPLAIN ESTIMATE PLAN:一个只读 visitor 遍历定稿的物理计划,把 CBO 已有的算子级估算吐成 JSON。纯加法(零删除),机制部分已准备为 upstream PR。 -
EXPLAIN ROUTE PLAN:估算加一个 15 规则分类器,6–8 毫秒返回路由裁决,与数据量无关。分类器从一个被当作可执行 oracle 的 Python 参考实现移植:任何集群介入之前先过 189/189 golden corpus parity。全部约 21 个阈值可热更,调参永不需要重新构建。
有一个判断和代码同样重要:leadership 起初想把分类器也贡献上游。我用第一个 PR 自己的原则反着论证:引擎应该吐数据,调用方决定策略。分类器携带我们校准的阈值和 recall 优先的偏置,那是策略,应该留在内部。知道什么不该贡献,改变了计划。
Heavy recall 在标准套件上从 7% 到 100%,剩余的漏网由设计兜住:误判的重查询在轻池撞上固定槽位内存硬限(约每查询 404 MB、禁止超卖),约 3 秒内被杀并自动升级到重池。分类器被允许不完美,因为安全保证从不依赖它。
两次多数运维者从未见过的 compaction 事故
**死亡螺旋。**迁移完成后,宽表冻结在每 tablet 约 2,500 个 segment;手动 compaction 静默失败;冷点查 40–68 秒。根因是三个因素连乘:bulk load 期间关了 auto-compaction,债务堆到单个任务要合并约 2,506 段 × 3,727 列(被内存守卫在 43 GB 处中止,任务级中止而非进程崩溃);cloud 模式没有 peer cache,每个 segment 都是约 7 秒的 S3 冷读;segment 之所以小,是因为稀疏宽 schema 在内存写缓冲和磁盘 segment 之间有 46 倍的压缩衰减。我们还在任何人花钱之前证明了加后端无用:单个 tablet 的 compaction 无法跨节点拆分,日志里数出 632 次空转唤醒对 9 次真实合并。是集中度受限,不是资源受限。验证过的配方(auto-compaction 常开、写缓冲扩大 4 倍、按月分桶把 tablet 钉在约 10 GB、加宽列组把 S3 重扫削约 40 倍)把冷点查从 40 秒压到服务端 100 毫秒以内。
**卡死的仪表。**第二次事故:compaction score 钉死在约 4,500,扩容和重启都无动于衷。三条独立的只读调查线(日志枚举、tablet 元数据 HTTP 端点、catalog 回收站)收敛到同一个孤儿 tablet:来自一张早已 DROP 的表,被结构性排除在调度之外。一条推理链会被带偏;三条独立证据链指向同一个 tablet 就是裁决。那次取证由并行 AI agent 分头扫过各个诊断面执行,这也是我现在跑取证的默认方式。
配得上「弹性」这个词的弹性
重池静息在零副本。scaler(一个我重建的内部 controller)只做一件事:对 Doris operator 的自定义资源 JSON-patch 副本数,纯 HPA 角色。它的第一版自己部署后端、用 SQL 注册;我在实现前就否掉了那个设计,因为 cloud 模式下任何没在 operator CR 里声明的东西都会被下一次 reconcile 抹掉。平台原语成熟时,自研工具就该收缩。
全链路真实跑通:裁决判重 → controller patch 0→N → spot 节点约 66 秒起来、后端约 2 分钟注册、分层就绪检查(SELECT 1 探针会被前端常量折叠、从不触达后端,它会把一个空池报告成健康)→ 查询执行 → 池回收归零。serving 池 live 从 2 台扩到 4 台,零数据搬迁:tablet 住在 S3,归属自动再平衡。这就是 shared-data 的运维红利,在真实集群上拍到的。
数字凭什么可信
上面每个论断都经过测量,其中几个在路上被撤回。早期 benchmark 显示 Doris 聚合赢 ClickHouse 3–5 倍;我 red-team 了自己的结果,找到混淆项(负载中的 8 核生产节点对隔离的 14 核 benchmark 节点),撤回了倍数。公开的「Doris 快 6–40 倍」全部溯源到厂商;一手 ClickBench 数据说 ClickHouse 仍领先平表扫描。迁移完整性认证为 99.945%(5,170,688,484 / 5,173,508,562 行),0.055% 的差额逐条归因到去重语义,而验证方法本身也得为安全重新设计:一次朴素的 COUNT(*) 对账曾打崩全部八台后端。
这一切之下的习惯:当一个数字每次重测都变得更不讨喜时,那通常说明你正在逼近真相。
文中的引擎级工作在一个 Apache Doris fork 上;机制(EXPLAIN ESTIMATE PLAN)已准备贡献上游,策略层刻意留在内部。名称、租户与精确标识已匿名化,规模数字取整。