Published on

从 OceanBase 到 StarRocks:基于 Kappa 架构构建现代实时数仓

Authors
  • avatar
    Name
    Charles Chen
    Twitter

随着实时风控、实时运营、实时推荐和实时监控成为常态,传统的 T+1 数仓已无法满足业务时效要求。

与此同时,企业的数据往往分散在多个数据库、网络、部门和机房中。难题不只是计算速度,更是如何让接入、运维和扩展保持可控。

本文以 OceanBase、CloudCanal、Kafka、Flink SQL 和 StarRocks 组成的一条实时链路为例,说明为什么它适合以 CDC 为主要数据来源的企业场景,以及各组件应承担的边界。

架构总览

整体采用 Kappa 架构,将业务变更持续转化为事件流,再由同一条实时链路完成加工和服务:

OceanBase
CloudCanal CDC
Kafka
Flink SQL
StarRocks

这条链路把数据采集、消息缓冲、实时计算和在线分析拆开。每个组件只处理自己擅长的工作,因此既能独立演进,也能在上下游维护或升级时保持系统韧性。

为什么选择 Kappa 架构

Kappa 架构的核心思想是:一切数据都是事件,一切计算围绕流进行。事件进入 Kafka 后,Flink 持续消费并写入服务层;需要修复或重算时,则从 Kafka 中按指定位置重新消费。

传统 Lambda 架构通常同时维护批处理和流处理两条链路:

Batch Layer ──► Serving Layer
Kafka ──► Speed Layer

两条链路意味着同一业务逻辑可能要分别由 Spark 和 Flink 实现。以订单 GMV 为例,无论离线还是实时,本质上都是 SUM(amount)

重复实现不仅抬高维护成本,也会让口径一致性更难保证。

Kappa 并不等于“永远不需要历史数据”。它的重点是让实时链路成为唯一的业务计算主线。

当历史重算、低成本存储或流批一体确实成为需求时,再增加合适的存储 Sink,而不是在需求出现前引入额外复杂度。

各组件的职责与选择

组件核心职责选择原因
OceanBase业务数据库(OLTP)分布式事务、高可用和高并发写入能力
CloudCanal数据采集层支持多数据源、多网段、CDC、断点恢复,并将采集与计算解耦
Kafka数据总线解耦上下游、缓冲削峰,并提供事件回放能力
Flink SQL实时计算引擎以一套 SQL 实现清洗、关联、聚合和指标计算
StarRocks实时分析数据库主键模型适合 CDC 更新,并提供高并发 OLAP 查询

OceanBase:只承担 OLTP

OceanBase 是业务事务的承载层,负责订单等实体的 INSERTUPDATEDELETE。它的目标是可靠地处理高并发事务,而不是承担大规模 BI 查询。

将分析压力直接压到业务库,会让事务负载与分析负载相互干扰。通过 CDC 将变更导出到实时链路后,业务库和分析系统才拥有清晰的资源边界。

CloudCanal:统一处理数据接入

很多团队首先会考虑让 Flink CDC 直接连接数据库。单一网络、少量数据源时这样做没有问题。

但当业务系统分布在不同网段时,Flink TaskManager 需要访问所有数据库,网络连通性、白名单和权限管理会迅速复杂化。

CloudCanal 将 CDC 独立为统一的数据采集平台,负责全量同步、增量同步、DDL 同步和断点恢复。Flink 只需面向 Kafka 消费,不再感知具体数据库的位置,从而实现采集与计算的解耦。

多个业务数据库
CloudCanal
Kafka

Kafka:为链路提供解耦与回放能力

Kafka 不是业务数据库,也不是计算引擎;它是整条实时链路的事件缓冲层。当 CloudCanal 暂停、Flink 升级或 StarRocks 维护时,消息仍可暂存在 Kafka 中,避免上下游必须同时可用。

更重要的是,Kafka 保留了可重新消费的事件日志。出现逻辑修复、目标表异常或需要补数时,可以通过重置消费位点或使用新的消费组进行 Replay,而不必再次从业务库全量抽取。

Flink 负责将原始变更加工为可分析的数据模型。常见的分层可以是:

  • ODS:字段清洗、标准化、脱敏和基础质量校验。
  • DWD:去重、关联和明细宽表构建。
  • DWS:窗口聚合与主题汇总。
  • ADS:面向报表、运营和业务指标的结果表。

使用 Flink SQL 的关键收益是统一表达方式。团队不必为批处理和流处理维护两套同义逻辑。

但是仍应为作业设计状态 TTL、Checkpoint、迟到数据策略、异常数据处理和幂等写入;否则,“实时”只会放大数据质量问题。

StarRocks:提供实时分析服务

CDC 场景中的订单状态会持续变化:创建、支付、取消和退款都会产生更新或删除。与主要面向追加写入的模型相比,StarRocks 的 Primary Key 模型更适合承接这类按主键更新的数据。

StarRocks 在这里是 Serving Layer:负责实时查询、高并发 BI 和 Dashboard,不负责业务事务,也不承担 ETL 的主逻辑。边界越清楚,系统越容易定位性能瓶颈和故障责任。

为什么暂时不引入 Paimon

Lakehouse 很有价值,但不应成为默认组件。当前链路的主要目标是实时分析,而不是 PB 级历史数据的低成本留存或长期重算。

在这个前提下,提前加入 Paimon 只会增加数据治理和运维复杂度。

当出现长期历史保存、低成本存储、离线训练或频繁流批重算的明确需求时,可以在 Flink 中增加新的 Sink:

Flink
 ├──► StarRocks:实时查询
 └──► Paimon:历史存储与重算

这种扩展不要求推倒现有链路,因为采集、计算和服务层已经被解耦。

为什么这里不选择 ClickHouse

ClickHouse 在日志、埋点和追加型 OLAP 场景中表现出色。但本文讨论的是以 CDC 为主的数据链路,存在大量按主键发生的 UPDATEDELETE

在这种场景下,StarRocks 的 Primary Key 模型通常更贴合业务实体的最新状态查询需求。技术选型不应只比较单点性能,而要看数据变更模式、查询模型和团队可维护性是否匹配。

企业落地时需要关注什么

这套架构解决的不是“哪个组件跑分更高”,而是如何让几十个业务系统稳定接入,如何降低网络与权限管理成本,以及如何统一计算口径。

同样重要的是,新增数据源时不应牵动整条链路。

落地时应至少确认以下事项:

  • 为每个 Topic 设计分区、保留期、Schema 演进和回放策略。
  • 为 Flink 作业明确 Checkpoint、状态 TTL、失败恢复、数据质量和告警标准。
  • 为 StarRocks 表模型确认主键、分桶、分区、写入幂等性和查询 SLA。
  • 为 CDC 链路演练全量初始化、增量衔接、DDL 变更、断点恢复和故障回放。
  • 对敏感字段在 ODS 或更早阶段完成脱敏,并记录数据权限边界。

总结

OceanBase 负责事务处理,CloudCanal 负责数据采集,Kafka 负责事件传输与回放,Flink 负责实时加工,StarRocks 负责在线分析。

每个组件只承担最适合自己的职责,才能形成高内聚、低耦合且可持续演进的数据平台。

真正优秀的架构,不是堆叠最多的新技术,而是在当前约束下,用最少的复杂度把职责边界划清楚,并为未来扩展留下稳定接口。