聊分布式系统的架构,最怕一上来就画一堆框框图。框框图大家都会画,真正难的是讲清楚每个框背后的取舍逻辑。这篇文章我打算认真聊聊 IDA Cluster,我们内部给它的全称是 Intelligent Data Architecture,智能数据架构,IDA Cluster 就是这套架构思想的落地实现。如果你在实时数仓、日志平台、数据中台这类方向工作,应该能看懂我在说什么。
先说一个真实现状:我经手维护过一条实时链路,从客户端埋点到最终业务看板,中间排布了 Nginx、Kafka、Flink、Redis、ClickHouse 五套系统。每套系统单独拿出来都很成熟,可是拼接起来就是另一个故事。某次大促压测,数据量一上来,Flink 反压导致 Kafka 消费 lag 飙升,ClickHouse 写入堆积,排障时三个同事同时盯着三个控制台,谁也说不清楚瓶颈到底在哪个环节。后来我们开始认真研究 IDA Cluster 的设计思路,核心就是解决这一类“组件缝合怪”问题:用一套集群把数据接入、轻量实时计算、结果存储和分发统一在一起。
这篇文章里,我不会写那种泛泛的架构科普。我会把我对 IDA Cluster 核心概念的理解拆开,包括节点角色、分区副本、任务状态、分层设计、与 Kafka/Flink 的边界、部署调优,以及我在真实压测中踩过的坑,一次性讲清楚。
1. 实时数据管道为什么会变成“组件缝合怪”
1.1 链路越长,出错的组合数越大
传统实时链路有一个固定的套路:数据源先进消息队列,再进流计算引擎,做一层处理之后落到各种存储里。消息队列用 Kafka,流计算用 Flink,存储用 ClickHouse、Redis、Elasticsearch,各有各的不可替代性。这种架构本身没问题,问题出在链条太长。
举个例子。一条订单支付事件,要经过客户端 SDK、网关、Kafka producer、Kafka broker、Flink source、Flink 窗口计算、Flink sink、ClickHouse 写入、最终报表查询。这中间任何一个环节抖动,都会往下游传导。更难受的是语义不统一:Kafka 自己有一套 at-least-once 的消费语义,Flink 靠 checkpoint 保证 exactly-once,ClickHouse 写入则依赖幂等表结构。三段拼起来,端到端到底是不是精确一次,没人能拍胸脯保证。
而且运维层面也很痛苦。Kafka 要盯 broker 的 ISR 收缩、磁盘水位、分区 leader 均衡;Flink 要盯 JobManager、TaskManager 内存、checkpoint 是否超时;ClickHouse 要盯 merge 任务和写入拒绝率。每个组件都有自己的监控大盘、告警规则和故障恢复方案,等于把一个完整问题拆成了三个不完整的问题,最后还得靠人肉去串联。
1.2 IDA Cluster 的设计初衷:把轻量计算下沉到数据接入层
IDA Cluster 的出发点,不是要取代所有组件,而是想把“从数据进来到结果可查”这条链路压缩到一套系统里。它给出的答案是:数据接入层不只是存消息,还要能直接做过滤、投影、窗口聚合、分发这些轻量运算;运算结果直接写入内建存储;整条数据流的元数据、任务状态、分区副本由同一个控制面管理。
具体来说,IDA Cluster 定位在三个“一体化”:
- 接入一体化:支持 HTTP、gRPC、TCP 等多种协议,数据进来之后做 schema 校验、格式解析、WAL 持久化。
- 计算一体化:内建轻量流式计算算子,可以用声明式配置描述数据处理逻辑,不需要单独部署一套 Flink。
- 存储一体化:热数据在内存或本地盘,温数据在本地 LSM 存储,冷数据可以转储到外部对象存储,保留策略由系统统一管理。
这套设计最直接的好处,是减少了组件数量。组件少了,端到端语义就更容易收敛。我在内部推动这个方案时,经常用一句话概括:如果一条数据进来只是做清洗、聚合、然后被查询,那它没有必要经过三套系统。
当然,IDA Cluster 并不是万能药。后面我会专门讲它的边界,以及什么时候不该用它。这里先明确一个大前提:它擅长的是中轻量级的数据处理,不适合超大规模状态计算和极其复杂的事件时间语义场景。
2. 核心概念拆解:节点角色、分区副本、有状态任务
2.1 四类节点角色,各自干什么活
一个 IDA Cluster 集群里,逻辑上存在四类节点,物理上可以混部。理解这四类角色的分工,基本上就理解了整套系统的骨架。
| 节点类型 | 核心职责 | 关键技术点 |
|---|---|---|
| 控制节点 Control Node | 元数据管理、选主、任务调度、权限控制 | 基于 Raft 的元数据复制,奇数节点部署 |
| 接入节点 Ingest Node | 接收外部数据、协议解析、WAL 写入、分区路由 | 批量刷盘、异步复制、背压感知 |
| 计算节点 Compute Node | 执行过滤、投影、窗口、聚合等任务算子 | 状态本地化、周期性 Checkpoint、watermark 推进 |
| 存储节点 Store Node | 保存原始数据与分析结果,处理查询请求 | LSM-Tree 结构、生命周期管理、索引 |
在小规模部署下,一台机器可以同时承担多种角色。比如三节点集群里,每个节点都跑 Ingest 和 Store,同时其中一台跑 Control,另外两台跑 Compute,甚至全部混部也能转。但无论怎么混部,Control 角色的实例数必须是奇数(通常是 1 或 3),因为 Raft 协议要求多数派才能选出 leader。如果只有两个 Control 实例,挂一个就无法形成多数派,集群控制面直接不可用。
我个人的建议是,生产环境低于五节点时,Control 用单实例加冷备就可以了,不必强行上三节点控制面;当集群规模超过十几台,再独立出三个 Control 节点。这套思路和很多存储系统的“轻控制面”设计是一致的。
2.2 分区与分片:逻辑上的分区决定了扩展上限
分区(Partition)和数据分片(Shard)是最容易被混淆的一对概念。分区是逻辑上的数据流切分单元,比如订单事件按 order_id 哈希分到 8 个分区,同一个订单的所有事件一定进入同一个分区,这样在分区内可以保证有序。分片则是物理存储上的承载单元,一个分区在每个副本节点上对应一个分片文件目录。
这里面的关键设计在于路由规则。客户端写入时,系统根据 key 做哈希映射到某个分区。虚拟桶哈希是常见做法:把 key 空间先映射到 1024 个虚拟桶,再让虚拟桶均匀分布到实际分区上。这样后续做分区扩容时,只需要重新映射部分虚拟桶,数据迁移量远小于直接按 key 范围切分。
分区数怎么定,一直是个经典问题。我的经验公式很简单:
预估峰值吞吐(MB/s) ÷ 单分区实测吞吐(MB/s) ≈ 分区数
单分区吞吐取决于磁盘性能和副本数,通常 NVMe 盘单分区可以跑 10~20MB/s。如果峰值是 100MB/s,那么 8~10 个分区是合理起步值。分区太多会放大元数据开销和文件句柄占用,分区太少又会限制并行度。宁可一开始略微偏多,也不要后续扩容时在线迁移数据——在线迁移的代价远比初始多几个分区要高。
另外一个细节是 offset 管理。每条消息在分区内有一个唯一序号,消费者任务需要记录已经处理到的位置。IDA Cluster 的控制面把消费位点存成一个内部元数据键值对,定期提交,任务恢复时从这个位点继续跑。位点提交有自动和手动两种模式:自动提交省事,但可能丢数据;手动提交更稳,但是代码复杂。我在内部任务里一律手动提交,尤其是下游还要写外部存储时,绝不能指望自动提交给你精确一次的保证。
2.3 副本数量与一致性:ack 级别不是越大越好
分区的每个副本里有且只有一个 leader。客户端写入只会发给 leader,leader 先写本地 WAL,再并行复制给 follower。副本确认数通过 ack 参数控制,这跟 Kafka 的语义很接近,但在实现上有个区别:Kafka 用 ISR(In-Sync Replicas)集合来动态判断哪些副本是同步的,IDA Cluster 基于 Raft 的日志复制机制,但允许业务数据走一种更高效的 quorum 确认模式。
三种 ack 级别,对应三种状态:
- ack=0:发送成功就算成功,可能丢数据,测试环境或者丢一点无所谓的指标采集可以用。
- ack=1:leader 写入 WAL 成功就返回,大多数业务场景的默认选择,延迟和持久性平衡最好。
- ack=all:所有同步副本都写入 WAL 才返回,持久性最强,但写入延迟显著上升。
很多人一上来就选 ack=all,觉得数据最重要,不能丢。这个想法没错,但你要为它付出代价:当某个副本因为磁盘繁忙落后,集群为了保证可用性会把该副本踢出同步集,此时 ack=all 实际上退化成 ack=当前同步副本数,而不是真正意义上的全部节点。反而因为频繁的副本状态切换,触发了不必要的告警。
我的做法是,核心交易链路用 ack=all + 同步副本数至少 2;普通日志链路用 ack=1 + 异步本地刷盘;测试链路用 ack=0。系统设计本来就是一个权衡,不是每个环节都需要银行级的可靠性。想清楚每条数据丢了会怎么样,你就知道该用哪档了。
2.4 有状态任务与状态快照:流式计算如何在集群里跑
接入层把数据写进分区之后,计算节点会启动任务去消费这些分区。一个任务在逻辑上是一张 DAG,节点是各种算子,边是数据流向。算子分两类:无状态算子和有状态算子。过滤、字段投影是无状态的,每个事件独立处理;窗口聚合、计数、去重则是有状态的,需要留存跨事件的状态数据。
状态不能只存在内存里,否则节点一挂全丢。IDA Cluster 的做法是周期性做 Checkpoint:把当前窗口数据和 keyed state 做一个分布式快照,保存到 Store 节点。当计算节点故障时,调度器会找一台新的机器,从最近的 Checkpoint 恢复状态,再回到消费位点继续跑。这个机制,和 Flink 的快照思想是一致的,只是简化了对对齐逻辑的要求。
为什么状态快照要落到 Store 节点而不是计算节点的本地盘?我当时也纠结过这个点,后来想明白了:状态快照的消费方不是原节点,而是未来可能替代它的任意节点。只有把快照放到一个独立于计算节点的存储层,才能做到故障后重新调度到别的机器也能加载。本地盘快照恢复快,但太过依赖节点存活;远端快照恢复慢一点,但是能真正实现“任意节点都能接手”。对 5~30 秒级恢复时间来说,把快照放在 Store 层是更稳妥的选择。
3. 架构分层:控制面与数据面分离,一条日志的完整旅程
3.1 控制平面:元数据、选主与任务调度
控制面是整个集群的“大脑”,它维护着所有主题、分区、任务、节点状态和用户权限。这部分数据量不大,但绝不能丢。所以它内部是一个使用 Raft 协议复制的小型 KV 存储,所有控制面的变更操作都通过 leader 来执行。
任务调度器也住在这里。它的职责是把计算任务合理地分配到 compute 节点上,分配策略有一个重要约束:数据局部性。一个任务如果在消费分区 P,调度器会优先把它分配到存储分区 P 副本的那台机器,或者跟这台机器网络距离最近的一台机器上。数据在本地读,比跨网络拉数据快一个量级。分布式计算里的 locality 原则,在这里同样生效。
另一个很实用的能力是任务版本管理。每一次任务配置更新,控制面会生成一个新版本号,然后以滚动方式逐步替换旧版本的计算任务。如果新版本运行异常,可以一键回滚到上一个版本。这种设计在接入层任务迭代频繁的场景下帮了大忙——我不用再靠人工记录“上次跑得好好的配置到底是什么”,系统自动帮我存了好几个版本。
3.2 数据平面:一条订单事件日志在集群里经历什么
为了把数据流讲得具体一点,我拿电商订单事件来走一遍完整链路:
- 业务服务端把一条订单事件以 JSON 形式 POST 到接入节点的 HTTP 接口。
- 接入节点解析 JSON,校验 schema 是否合法,然后按 order_id 哈希到某个分区,写入 WAL,返回成功给客户端。
- WAL 刷盘之后,复制线程把日志推送给该分区的 follower 副本,副本写入成功后更新同步进度。
- 计算节点上运行着订单实时指标任务,从这些分区消费数据。任务先做一次过滤,把非订单事件丢弃;再做字段投影,只保留 order_id、user_id、amount、status、ts 等必要字段。
- 窗口算子按照事件时间把数据切分成 60 秒一个窗口,窗口结束时触发聚合计算,产出这一分钟的订单数、GMV、支付成功率。
- 聚合结果写入 Store 节点下的一张结果表;原始日志按照保留期限,在 7 天后被自动清理。
这条链路全部发生在一个集群内部,不需要跨 Kafka、Flink、ClickHouse。它有另一个隐蔽的好处:跨组件调试成本急剧下降。以前查一条数据丢了,要翻三个组件的日志;现在只要在控制面按 traceId 查一次流转记录,就能定位到具体是接入、计算还是存储哪个环节出的问题。
3.3 异步刷盘与背压:为什么不能靠缓存解决问题
分布式链路里,每个环节处理速度不一样,就会出现上下游速度不匹配。很多初学设计的人第一反应是“加缓存、加队列”,让快的先积压着,等慢的慢慢消费。这听起来合理,实际上是个陷阱:内存是有限的,积压到一定程度必然要淘汰数据或触发 GC,结果就是延迟毛刺和丢数据一起出现。
IDA Cluster 采用 credit 制流控,核心思路很直白:消费方明确告知生产方自己还能接收多少数据,生产方在 credit 用完之后就必须停下,等待新的 credit。
你可以把它理解成两个人搬砖:楼上的人只告诉楼下的人“我还能再接一箱”,楼下的人绝不会一次性把十箱都堆在楼梯口。credit 制的好处是,背压能一级一级传导回去。如果 Store 节点写入慢了,Compute 节点就停下来不生产;Compute 停下来,Ingest 节点的数据就开始在 WAL 里堆积;WAL 堆积到阈值,接入节点开始拒绝新的外部写入,同时返回明确的“服务过载”错误码。
这套机制保证了任何环节的瓶颈都会真实地暴露出来,而不是被内存缓冲区掩盖。真实排障中,我看到太多案例都是“消息队列消费 lag 很高,但组件本身内存还够”的假象,本质上就是因为某个下游静默变慢,而上游毫无感知。背压存在的一个价值,就是让慢的环节无处可藏。
4. 和 Kafka、Flink 划清边界:选型不是越多越好
4.1 IDA Cluster 不是 Kafka 的替代品
我最早看 IDA Cluster 的文档时,心里也在嘀咕:这个东西有主题、有分区、有消费位点,不就是 Kafka 吗?后来深入用才想明白,它俩关注的根本不是同一个层次的问题。
Kafka 是一个分发系统,它关心消息怎么持久化、怎么被消费者拉取、怎么在 consumer group 之间做负载均衡。它不关心消息内容是什么,更不关心消息接下来要算什么。IDA Cluster 除了分发,还关心消息进入之后该怎么解析、该触发哪些任务、结果该写到哪个结果表,它是一个面向“数据加工结果”的系统。
| 维度 | Kafka | IDA Cluster |
|---|---|---|
| 核心抽象 | 分区日志 | 数据架构:接入 + 任务 + 存储 |
| 消息保存 | 按保留期存储,消费后不删除 | 同样按保留期存储,但多了任务视图 |
| 内建计算 | 无 | 有轻量流式计算算子 |
| 消费模式 | Consumer Group 重平衡 | 任务订阅 + 即席查询 |
| 端到端管理 | 不管数据从哪来、到哪去 | 管理从接入到存储的全链路 |
所以更准确的说法是:Kafka 适合做总线型基础设施,IDA Cluster 适合做端到端的数据产品底座。如果你的团队已经有了成熟的 Kafka 基础设施,完全可以把 Kafka 放在接入前端,IDA Cluster 消费 Kafka 再做聚合和存储;如果你是从零开始建一套实时数据平台,不想维护那么多组件,直接上 IDA Cluster 会省心很多。
4.2 与 Flink 的差异化与协作模式
Flink 依然是流计算领域的事实标准,它的算子生态、状态管理、事件时间语义都非常成熟。IDA Cluster 内建的算子,更适合清洗、聚合、计数这类确定性很强的轻量任务,而不是复杂的业务规则计算、CEP 复杂事件处理或大规模机器学习特征计算。
我的实际分工原则是:
- 默认尽量在 IDA Cluster 内完成。过滤、字段映射、1 分钟窗口聚合、滚动指标,这些是数据架构中最常见的操作,用声明式配置就能搞定,没必要额外铺一套 Flink。
- 超出轻量边界的任务外派给 Flink。状态规模预计超过单节点内存、需要自定义 UDF 与外部系统深度集成、要精确处理乱序事件并做多级窗口 join,这些情况就让 IDA Cluster 把清洗后的干净数据交给 Flink,由 Flink 完成重型计算。
- 不要让两层做重复的数据清洗。如果 IDA 已经把字段规范化了,Flink 就不需要再解析一遍原始 JSON,直接消费规范后的 schema,这样能省掉大量无谓 CPU。
我之前见过一个项目,数据从 Kafka 进 Flink,Flink 里先做一层 JSON 解析和脏数据过滤,然后下游的 ClickHouse 表还要求数据再做一遍类型转换。三层各做一遍同类工作,性能和维护成本都是灾难。如果中间由 IDA 做统一的 schema 校验和清洗层,Flink 只做最核心的计算,事情会简单得多。
4.3 什么场景不该用 IDA Cluster
任何技术都有适配边界。我在团队内部反复强调,不要为了统一而统一,下面几种场景就不建议强行用 IDA Cluster:
第一,已有的 Kafka + Flink 体系已经跑得很稳,且监控、告警、运维流程都很成熟。迁移成本大于收益,这时候再造一套反而破坏稳定性。
第二,业务涉及超大状态或者复杂的时间语义。比如按用户维度保留 90 天行为序列,每天的状态规模上百 GB,这种场景需要 Flink 的 RocksDB 状态后端和增量 Checkpoint 配合,IDA Cluster 轻量状态架构撑不住。
第三,团队只需要一个消息管道。如果需求就是把 A 系统的数据搬到 B 系统,不做任何加工,那直接用 Kafka 会轻得多,把计算和存储能力都引进来反而增加心智负担。
还有一点容易被忽略:架构统一不等于团队技能统一。引入 IDA Cluster,意味着团队要熟悉一套新的配置语法和运维工具。上线前一定要预留学习和试错的时间,不要指望三天内所有业务线都能切换上去。
5. 部署与调优:三节点集群的实战经验和踩坑记录
5.1 三节点最小化部署:角色怎么分配
我建议的最小生产集群是 3 台机器,规格 16 核 32G 内存起步,本地磁盘最好用 NVMe SSD。低配机器虽然能跑通,但一旦开启多任务和窗口聚合,CPU 和内存都会很紧张。
三台机器角色分配可以这样:每台都跑 Ingest 和 Store,因为接入和存储是数据密集型的,三节点天然分散压力;Control 角色保持单 leader 加两个 follower 部署在三台机器上,组成一个 Raft 组;Compute 任务则通过调度器自动分配到当前负载最低、且离数据最近的节点。
这里有个容易犯的错:以为 Control 节点越多越稳。实际上 Raft 组三节点和单节点相比,每一次元数据变更都需要多数派确认,写入延迟会增加。小集群里控制面变更频率并不高,真正的瓶颈从来不在控制面。日常运维中,我甚至建议把控制面的 leader 固定在与外部客户端网络质量最好的那台机器上,避免 leader 频繁漂移带来额外抖动。
5.2 关键参数怎么看:从分区数到 Checkpoint 间隔
部署配置里最核心的几个参数,我直接给一组经过压测参考值:
| 参数 | 推荐值 | 含义与调整建议 |
|---|---|---|
| partition count | 8 | 按峰值吞吐计算,宁可略多不可太少 |
| replication factor | 2 | 3 节点下容忍单节点故障,重副本因子 3 会拖慢写入 |
| flush interval | 512 条 或 200ms | 条数和时间满足其一就刷盘,适合中等延迟敏感业务 |
| checkpoint interval | 30s | 故障恢复越快则间隔越短,但频繁快照会占磁盘 IO |
| compute parallelism | 等于分区总数 | 一个分区同时只被一个计算线程消费,保证分区内有序 |
| store memory budget | 节点内存的 50% | 给 LSM 的 block cache 分配,别让 Compaction 抢走所有 IO |
关于副本因子,我想多说一句。三节点集群很多人会配副本因子 3,觉得数据冗余越足越安全。但副本文本不是白来的:每次写入都要复制两份,网络带宽和磁盘占用翻倍。对一个三节点集群来说,坏两台机器的概率非常低,副本因子 2 已经能在单节点故障时保住数据。副本因子 3 更适合节点规模超过 5 的集群,那时候单节点故障期间还要再坏一台的概率才真正值得用第三副本去对冲。
5.3 压测中发现的问题,比文档更有价值
我把真实压测中遇到过的几个问题列出来,这些问题在官方文档里大概率看不到,但实践价值很高。
| 现象 | 根因 | 解决方式 |
|---|---|---|
| 某个分区消费 lag 持续走高,其他分区正常 | key 分布不均,某个大客户订单量占比超过 30% | 写入 key 加盐,或者采用两层分区:先按租户分流,再按订单哈希 |
| 磁盘 IO 到达瓶颈,但 CPU 很闲 | 窗口聚合的中间结果频繁刷盘,涉及太多分片文件 | 调大 Store 的 memtable 阈值,减少小文件 Compaction |
| 写入 P95 延迟突然从 100ms 涨到 1s | batch.size 设置太大,刷盘等待时间过长 | 调小批次大小,同时在接入层增加 in-flight 请求限制 |
| 任务重启后恢复时间长达几分钟 | Checkpoint 太频繁,快照文件过多 | 调大 checkpoint 间隔,同时开启增量快照 |
压测最值得注意的一个教训是:不要只看平均延迟。实时系统里,P99 和 P95 才是用户真实感受的上限。某些问题会让平均延迟只涨 30ms,但 P99 已经翻了几倍。所以压测报告里,我会同时盯平均值、P95、P99 和“最大分区 lag”四个指标,任何一个异常都不能放过。
5.4 故障排查思路:从现象到根因的完整链路
这里分享一次真实的排障过程。某天线上监控告警:“接入节点拒绝写入”,我第一反应是磁盘满了。登录机器看了下磁盘,剩余 30%,并没有满。
继续查 WAL 写入日志,发现 WAL 目录的写入耗时从平时的 2ms 涨到了 200ms。这才意识到问题不在接入节点本身,而是它的 WAL 文件所在的磁盘卷和 Store 节点的 Compaction 任务共享了同一块磁盘 IO。大范围 Compaction 把磁盘 IO 打满,WAL 写入跟着变慢,堆积到阈值后接入节点开始拒写。
解决办法分两步:第一步,把 Ingest 的 WAL 数据目录和 Store 数据目录放到不同的磁盘挂载点上,从物理上隔离;第二步,给 Compaction 任务加了 IOPS 上限,避免它占用全部磁盘带宽。从那以后,我再也没有遇到过磁盘争抢导致的接入拒写。
这个案例想说明一个排查思路:遇到表面现象不要急着在处理层找原因,往时间线上游多看一层。接入拒写,问题可能在存储;计算延迟,问题可能在副本复制;任务反复重启,问题可能在状态快照恢复。分布式系统里的故障,十有八九不在第一现场。
6. 用一套订单实时监控把全部概念串起来
6.1 场景与核心指标
假设现在要做一个电商订单实时看板,业务方要求的指标是:每分钟的订单数、GMV、支付成功率、支付 P95 耗时。数据源是服务端各业务系统上报的 order_event 事件,里面包含 order_id、user_id、amount、status、pay_cost_ms、event_time 等字段。订单状态有 created、paid、cancelled 三种,支付成功率就是 paid 数量除以 created 数量。
这个场景非常适合 IDA Cluster:数据量中等(峰值每秒几万条),计算逻辑简单(过滤、窗口聚合),结果需要实时查询。换成 Kafka + Flink + Redis 或 ClickHouse 的架构也能做,但要管理和运维的组件就多出好几套了。
6.2 建表与任务配置
下面是一份可参考的配置。注意这不是完整的生产配置,只列出核心字段,帮助你建立直觉。
source: type: http path: /v2/order_events protocol: grpc table: name: order_event_raw partition: 8 replication: 2 retention: 7d task: name: order_realtime_metrics from: order_event_raw steps: - filter: event_type == "order" - project: order_id, user_id, amount, status, pay_cost_ms, event_time - window: tumbling(60s, event_time) - aggregate: group: [window_end] create: order_count select: - count(*) as order_cnt - sum(amount) as gmv - sum(if(status == "paid", 1, 0)) / count(*) as pay_rate - percentile(pay_cost_ms, 0.95) as pay_p95 sink: type: store table: order_metrics_1m这段配置表达的意思是:从 order_event_raw 读数据,先过滤,再做 60 秒的滚动窗口聚合,结果写入 order_metrics_1m 结果表。filter 和 project 是无状态算子,window 和 aggregate 会触发状态快照。整个任务的逻辑清晰可见,版本更新也只是改配置然后发布新版本,不需要写 Java 或者 Scala 代码。
6.3 压测结果与问题分析
在三节点集群、8 分区、2 副本的配置下,我做过一轮压测。用压测工具模拟客户端上报,数据从 1 万条/秒逐步加到 20 万条/秒,观察各阶段的延迟和系统资源变化。
| 数据量 | CPU 使用率 | 磁盘 IO | P95 写入延迟 | 结果表查询 P95 |
|---|---|---|---|---|
| 2 万条/秒 | 35% | 25% | 180ms | 40ms |
| 8 万条/秒 | 60% | 55% | 280ms | 90ms |
| 15 万条/秒 | 82% | 90% | 520ms | 230ms |
| 20 万条/秒 | 95% | 97% | 超过 1s | 超过 1s |
可以看出,当磁盘 IO 达到 90% 以上,写入延迟急剧恶化。此时 Store 节点成为瓶颈。我把结果表 order_metrics_1m 的分片数从 8 调整为 16,让 Compaction 的粒度变小、并行度提高,20 万条/秒场景的 P95 写延迟降到了 650ms。如果继续加大分区数可能还能再压一点,但收益已经开始递减。
这个结果说明一个原则:流式系统压测,瓶颈往往不在计算而在存储层的写放大。量化任务时,第一先估算写入和存储的带宽,再去调计算并行度。
6.4 重新设计这套系统时,我会重点改的三个地方
第一,我会给订单事件加上多级分区。直接按 order_id 哈希,头部大客户的订单会把几个分区打胖。更好的做法是先按租户或业务线做一层分流,再在分片内部按 order_id 哈希。这样单个分区的数据热度更均匀,也不用担心某个大客户拖垮整个集群。
第二,我会把原始日志的冷数据尽早转储到对象存储。本地盘存 7 天原始数据没问题,但如果要存 30 天甚至 90 天,磁盘成本太高。冷热分离应该从一开始就设计进去,而不是等磁盘告警了再补方案。IDA Cluster 支持将超过保留期限的分片自动归档到 S3 兼容存储,这个能力值得优先使用。
第三,我会在任务配置里增加一个“查询结果缓存”层。看板场景的特点是同一份指标会被不同页面频繁查询,如果每次都打 Store,查询压力很容易成为新的瓶颈。一个简单的做法是把最近 5 分钟的聚合结果缓存到控制节点的内存里,过期后自动失效,查询 P95 能下降一个数量级。
我个人的实际体会是:在实时数据平台的选型上,组件数量真的不是越多越有底气,而是每一层都要回答清楚“它到底承担了什么不可替代的职责”。IDA Cluster 这套思路,好就好在把数据接入、轻量计算和结果存储收敛到一起,让我可以把精力放在业务指标和数据处理逻辑上,而不是整天盯着各个组件之间那根细若游丝的链路。
如果你也在为一条动不动就从 Kafka 经过 Flink 再到 ClickHouse 的链路头疼,不妨先拿一个订单监控或日志清洗这种中等场景,试试把链路缩短。你会发现,很多所谓“架构问题”,其实只是组件摆得太多而已。