简介:OPPO基于Apache Flink的实时数仓实践演示文稿,面向大数据工程师、数仓架构师及流计算开发者,系统阐述以Flink为核心构建秒级响应的实时数仓方案。内容围绕背景、高级设计、实践架构、最佳实践与未来工作展开,结合ColorOS数亿月活用户的业务场景,详细讲解离线与实时数仓的平滑迁移、ODS/DWD/ADS分层建模、统一数据采集管道、Flink SQL开发与元数据管理、Kafka主题重复消费优化、数据链路自动化调度等落地要点,并给出实时报表、实时用户画像等典型应用的工程经验,为同类业务场景提供参考路径。资料仅含一个pptx文件,大小10.39MB,共四大部分内容层次分明,适合作为团队内部分享或方案设计参考。该演示文稿已有407人学习浏览,值得正在规划或优化实时数仓的技术团队研读借鉴。
1. 实时数仓这条路,OPPO踩过的坑不比任何人少
Apache Flink 做实时数仓,在 OPPO 手里不是 Demo,而是扛住了 ColorOS 三亿月活用户、几十个应用(主题商店、浏览器、游戏中心、应用商店、云服务、搜索、小游戏)的实时报表、实时画像和实时接口。这套实战经验沉淀下来形成的 PPT,核心就一句话:实时数仓并不是另起炉灶,而是和离线数仓共用一套采集、一套元数据、一套权限体系,只是在时效性上从小时/天级压到秒级。适合谁?正在从离线数仓向实时数仓演进、被 Kafka topic 重复消费和数据延迟折磨的团队;也适合想把 Flink SQL 真正用进生产、而不只是跑通 WordCount 的开发。下文我把这份实践拆成可复现的思路,包括分层设计、SQL 统一、JobGraph 优化和避坑细节。
2. 背景与高层设计:从离线准时跑到实时秒级,OPPO为什么选 Flink
2.1 业务压力:三亿月活背后的实时诉求
OPPO 的移动互联网业务覆盖了主题商店、浏览器、游戏中心、应用商店、云服务、搜索和小游戏等十几个应用,ColorOS 月活跃用户数超过三亿。这个体量带来的数据需求分三类:第一类是实时报表,比如广告的曝光量、点击率,业务方要看到分钟级甚至秒级的变化;第二类是实时画像,比如用户当前所在位置,要能基于最新一条位置记录更新;第三类是实时接口,比如用户下载某个 App 后,云端服务要立刻感知并触发后续逻辑。
这三类需求在离线数仓时代都做不了。原来的架构是典型的数据仓库中心式:Flume/NiFi 做数据采集,Hive/Spark 做批处理,Presto/Hive 做交互式查询,结果落到 Elasticsearch、MySQL、Kylin、Redis、HBase 等存储里,供 BI 报表、用户画像和服务接口使用。这个架构的瓶颈很明确:ETL 任务基本都集中在凌晨跑,对集群压力极大,而且产出的数据是小时级或天级的,业务方看实时报表只能干等。
我摘录 PPT 里的一个关键对比:离线数仓和实时数仓在数据源、数据分析师、数据应用层面是相似的,差异点主要在时间敏感度。离线数仓按小时/天产出,实时数仓按分钟/秒产出。OPPO 的策略不是把离线体系推翻重来,而是「平滑迁移」——让实时和离线共用数据源、SQL 语义、UDF、权限和血缘追踪,只是在存储引擎和执行引擎上做了区分:Hive 管离线,Flink 管实时。
2.2 统一采集管道:OBus 与 Kafka
PPT 里最大的一张架构图是统一采集管道:OBus 接入业务数据,写入 Kafka,再由 Kafka 分流到两个方向——HDFS 走离线链路,Flink SQL 走实时链路。这个设计看起来简单,但它是整个实时数仓的地基。
这里的核心思想是「一份数据入口,两条处理路径」。采集端只做一次,不管下游是离线还是实时,都从 Kafka 里拿数据。这样做有几个实际好处:一是采集端的部署和运维只有一套,不用为实时单独建采集任务;二是 Kafka 本身充当了削峰填谷的缓冲层,凌晨的波峰不会直接打垮计算引擎;三是离线链路和实时链路可以互相校验数据一致性,数据源相同,结果不应有本质冲突。
OBus 是 OPPO 自研的数据总线组件。它负责从业务数据库、日志文件、消息队列等来源拉取数据,并保证至少一次投递到 Kafka。我在实际项目里的经验是,采集端最容易出问题的不是并发不够,而是数据格式不统一——有的业务方给 JSON,有的给 Avro,有的直接给 CSV。OBus 的做法是在源头做了一次归一化,统一成 Kafka 里的标准消息格式,这样下游 Flink SQL 解析时不需要为每个业务都写一套 deserializer。
提示:如果你没有 OBus 这类自研组件,用 Canal 同步 MySQL Binlog 到 Kafka,或者用 Flume/Kafka Connect 采集日志,同样能构建统一采集管道。重点是采集层只做一次,不要让实时和离线各搭一套。
2.3 统一管理进程:元数据、权限、监控与血缘
PPT 里有一页专门讲统一管理进程,这是容易被忽略但真正决定实时数仓能否长期跑稳的部分。它包含四块:元数据系统、权限系统、监控系统、表血缘追踪。
元数据系统负责管理「表」的定义。在 OPPO 的架构里,Hive 表和 Flink 表共用同一套元数据,Flink 通过 ExternalCatalog 把 Hive 的元数据转换成 Flink 的 Table 定义。这样用户写 SQL 时不用关心这张表到底是离线表还是实时表,语义是一致的。
权限系统延续离线数仓的权限体系,实时表同样纳入行列权管控。监控系统覆盖作业级和 topic 级,作业挂掉、消费延迟、checkpoint 失败都要能告警。血缘追踪用来回答「这张表的数据是从哪张源头表来的、被哪些下游任务消费了」。这些能力在离线数仓里相对成熟,难点是把它们延伸到实时场景,因为 Flink 作业是常驻运行的,血缘关系是动态变化的。
我做实时数仓的经验是,很多团队一开始只关注计算逻辑,把元数据、权限、血缘这些「管理能力」往后放,结果数据链路一长就失控——某个 topic 被改了 schema,下游十几个作业静默报错;某个任务出了问题,找不到上游是谁产生的脏数据。OPPO 这个「统一管理」的思路值得借鉴:在搭建实时数仓的初期就把管理和计算并行建设。
3. 数仓分层与核心组件:ODS/DWD/ADS 在 Kafka 和 Flink SQL 里怎么落地
3.1 分层设计:每一层用什么引擎,为什么
PPT 把实时数仓分成 ODS、DWD、ADS 三层,外加 DIM 维度层。这个分层模型是离线数仓经典分层在实时场景的映射,但每一层的实现工具差异很大。
ODS(操作数据层)直接对接 Kafka,数据从 OBus 进来后,Flink SQL 只做清洗、过滤、格式转换,不join不聚合,然后写回 Kafka 或者 HDFS。这一层的目标是「原样保留,去重去噪」。DWD(明细数据层)做核心的加工逻辑——join、维度补充、事件整理、会话拆分,产出干净的明细事实,存储介质仍是 Kafka。ADS(应用数据层)面向具体业务需求,比如曝光点击率的分钟级聚合,产出结果可能落到 Kafka,也可能直接写入 Druid、Elasticsearch 等查询引擎。
DIM 维度层存放在 MySQL 或者 Hive 中,Flink SQL 在 DWD 层通过维表 join 补齐维度信息。PPT 里的架构图显示:
ODS: OBus -> Kafka -> Flink SQL -> Kafka DWD: Kafka -> Flink SQL -> Kafka ADS: Kafka -> Flink SQL -> Kafka -> Druid/ES DIM: MySQL/Hive -> Flink SQL 维表关联这个链路里有一个容易被忽略的设计决策:ODS 和 DWD 层为什么都用 Kafka 作为存储?因为在实时链路里,Kafka 承担了「消息中转+数据存储」的双重角色,上下游作业通过 topic 解耦。DWD 层加工完的结果写到一个新的 topic,ADS 层的作业消费这个 topic 做聚合。如果 DWD 的结果还要供离线使用,可以再加一个 sink 写到 HDFS。按需双写,而不是所有数据都双写。
3.2 一套 SQL 语言:批量、流式、交互式统一
PPT 里有一页配图非常有趣,标题叫「One SQL to rule them all」。中心是 Stream 数据仓库,周围是 Query 系统、Ingestion 系统、Service、Development 系统、Business 系统,全部以 SQL 为交互语言。这页图表达的是 Flink SQL 在 OPPO 的角色不只是计算引擎,而是整个实时数据生态的「通用语言」。
传统做法是每个环节用不同的工具和语言:采集用 Flume,处理用 Java 写 Storm/Spark Streaming,查询用 Presto,开发平台自研一套配置。OPPO 的做法是把这些统一到 Flink SQL 上。收益很明显:学习成本低,离线数仓的开发者能平滑迁移;逻辑复用度高,同一个 UDF 在离线和实时场景都能用;排查问题方便,一条 SQL 的逻辑一眼能看穿。
当然,一条 SQL 很方便意味着 SQL 本身能力要足够强。Flink SQL 在 OPPO 的实践里必然遇到了这些问题:复杂的维表 join 怎么异步查询、状态清理怎么控制 TTL、CDC 数据如何处理更新删除。这些 PPT 没展开,但按 Apache Flink 实际社区发展和 OPPO 当时的落地时间点(PPT 里出现的是较早版本),核心是「用 SQL 描述流处理逻辑,让引擎去处理时间和状态」。
我在生产里用 Flink SQL 的体会是:能用 SQL 写清楚的就别写 DataStream API。SQL 让流处理的复杂性(时间窗口、watermark、状态管理)隐藏在引擎内部,业务代码量可以压缩一个数量级。代价是调试不够直观,所以需要配套的作业拓扑可视化——OPPO 用的是 AthenaX 平台,我下文会展开。
3.3 元数据管理:Flink ExternalCatalog 的表生命周期
PPT 里专门有一页讲元数据系统的实现流程,流程大致如下:
MySQL 元数据库 -> createTable (创建表定义) -> Flink ExternalCatalog (转换成 Flink 能识别的目录结构) -> convert to Flink Table / register (注册为 Flink 临时表) -> persist / load (持久化到元数据库)这条链路解决的是一个实际问题:Flink 作业提交前,表定义是存在 Flink 内部的,作业停止后表定义就丢了;而离线数仓的表定义是存在 Hive Metastore 里的,长期有效。OPPO 的做法是在 Flink 和元数据库之间架了一个 ExternalCatalog 适配层。
这个 ExternalCatalog 帮我理解 OPPO 的实时数仓研发流程:开发者在 AthenaX 里建表,元数据写入 MySQL;Flink 作业启动时,从 MySQL 加载表定义,注册成 Flink 的 Table 对象;作业运行中用到的表如果第一次出现,则动态创建对应的 Kafka topic 或 HDFS 目录。这样表和实际物理存储的生命周期是统一管理的,不会出现「SQL 里引用了表但不知道对应哪个 topic」的混乱。
我在项目里做类似设计时的做法是用 Hive Metastore 作为统一元数据中心,Flink 通过 HiveCatalog 对接。OPPO 选择自研 ExternalCatalog 对接 MySQL,我猜是为了能同时管理 Kafka topic、HDFS 目录、Druid 数据源等多种物理存储类型。这个设计本身说明一个道理:实时数仓的元数据管理要比离线数仓更「多模态」,因为下游存储引擎太多了。
4. 开发与调度:从建表、提交作业到 YARN 上的 JobGraph
4.1 AthenaX:SQL 开发与作业提交链路
PPT 里有一页详细展示了 OPPO 的 SQL 开发和提交链路。开发系统(AthenaX)提交作业的流程大致如下:
开发者编写 SQL -> AthenaX (SQL 平台) -> JobStore (作业持久化) -> FlinkTableEnvironment (编译 SQL 为执行计划) -> compile (生成 JobGraph) -> YARN Client (提交到 YARN) -> YARN 集群运行 Flink JobAthenaX 是 OPPO 自研的 SQL 开发平台,它的职责不只是「写 SQL」,还包括:语法检查、作业版本管理、权限审批、作业调度、运行状态监控。PPT 里提到了 JobStore 组件,说明作业定义是持久化的,不是一次性的提交任务。
开发者在这个平台上的典型操作序列是:第一步注册表(source/sink 表),选择数据源类型(Kafka topic 或 HDFS 路径)、消息格式(JSON/Avro)、字段映射。第二步写 SQL 逻辑,在平台上做语法检查和逻辑校验。第三步配置作业参数,如并行度、checkpoint 间隔、状态后端、重启策略。第四步发布作业,平台提交到 YARN,然后持续监控。第五步如果需要更新逻辑,提交新版本或回滚。
这条链路的可复现性很强。即使你不用 AthenaX,用 Flink SQL Gateway 或开源的 StreamPark/Dinky 也能搭出类似闭环。核心是「作业定义版本化 + 提交过程自动化 + 运行状态可观测」。
4.2 作业提交的参数与调优:并行度、Checkpoint 与状态后端
实时数仓里,Flink SQL 作业的参数设置直接影响稳定性和延迟。OPPO 的 PPT 没有逐个列参数,但基于 Flink 生产实践的通用做法,我会按这套逻辑来配置:
# 伪代码:Flink SQL 作业的推荐参数配置(以 Flink 1.13+ 为例) env.set_stream_time_characteristic(TimeCharacteristic.EventTime) # 事件时间 table_env.get_config().set_local_time_zone("Asia/Shanghai") # 时区 # 核心参数 parallelism.default: 4 # 按 Kafka topic 分区数的一半到一倍设置 state.backend: rocksdb # 大状态用 RocksDB,避免堆内存溢出 state.checkpoints.dir: hdfs://nameservice/flink/checkpoints # HDFS 存 checkpoint execution.checkpointing.interval: 60s # 60 秒一次 checkpoint execution.checkpointing.tolerate_failed_checkpoints: true # 一次失败不致命 execution.checkpointing.min-pause: 30s # 两次 checkpoint 之间最小间隔 restart-strategy: fixed-delay # 固定延迟重启 restart-strategy.fixed-delay.attempts: 3 # 最多重试 3 次 restart-strategy.fixed-delay.delay: 10s # 重启间隔 10 秒 # Kafka 参数 connector.properties.group.id: flink_sql_group connector.properties.auto.offset.reset: earliest scan.startup.mode: group-offsets # 从 group 提交位点开始消费并行度设置有一个常见教训:不是越大越好。Flink SQL 作业的并行度直接决定 Kafka 分区被怎么分配、状态被怎么切分。我一般按 Kafka topic 总分区数设并行度,上限不超过分区数的两倍。并行度超过分区数时,多出的 Subtask 空转浪费资源;低于分区数时,一个 Subtask 消费多个分区会增大单点压力。
Checkpoint 间隔需要根据「状态大小」和「业务容忍的恢复时间」来权衡。状态几百 GB 的作业,每 60 秒做一次 checkpoint 会造成较高的 IO 开销;状态较小的作业可以把间隔压到 30 秒。OPPO 服务几亿用户,实时画像类作业的状态一定不小,用 RocksDB 是正确选择——避免状态增长撑爆堆内存。
4.3 实时管道的自动化:Kafka 表到 BI 系统的数据流转
PPT 最后部分讲了工作流自动化,链路如下:
Kafka Table -> Flink SQL -> Kafka Table -> Druid -> BI 系统 数据加工 流式计算 结果缓存/引擎 可视化展示这条链路对应的是实时数仓从数据处理到业务可视化的完整闭环。Flink SQL 作业消费上游 Kafka topic,计算结果写到下游 Kafka topic;再通过 Druid 的实时导入能力把数据索引化,供 BI 查询。
我这里补一刀:Druid 在实时数仓里是个很独特的角色。它支持实时导入 Kafka 数据并自动构建索引,查询侧又是亚秒级响应,特别适合「数据不断进、查询随时出」的实时报表场景。OPPO 把 ADS 层的结果放 Kafka,让 Druid 消费 Kafka 构建索引,而不是 Flink 直接写 Druid,这样解耦了计算和存储,也给 druid 自身留了缓冲余地。
如果不用 Druid,用 ClickHouse 替代也能实现类似效果——Flink SQL 把结果写 Kafka,再通过 ClickHouse 的 Kafka Engine 表消费并写入本地 MergeTree 表。关键点是只让最下游的 ADS 层对接 OLAP 引擎,不要让 ODS/DWD 层直接写 OLAP 引擎,否则 OLAP 引擎既要扛写入压力又要扛查询压力,很容易被打爆。
5. 避坑指南:Kafka 重复消费、数据倾斜、延迟波动与状态恢复
5.1 同一个 Kafka topic 被多个 SQL 重复消费
现象:同一个作业里写了多条 SQL,都读取同一个 Kafka topic 作为 source,作业运行后 Kafka 侧的消费总量变成原来几倍,topic 的分区吞吐急剧上升。
原因:Flink SQL 在生成 StreamGraph 时,每个 SQL 引用同一个 topic,默认会生成独立的 DataSource 节点,相当于对同一个 topic 启动了多份消费者。
解决:OPPO 的做法是重写 StreamGraph,找出消费相同 topic、相同 group 的 DataSource 节点,合并成一个节点,让所有下游分支共享这一份消费。用 Flink SQL 提交作业时,可以通过 TableEnvironment 的解释器拿到 StreamGraph,再做 DataSource 节点的去重合并。这个操作要放在 Flink SQL 编译阶段,在提交 YARN 之前完成。
原理说明:StreamGraph 合并的本质是数据流复用。多个 SQL 都要消费同一个 Kafka topic,但每条 SQL 的投影和过滤条件不同,可以在「共享 DataSource → 分别做 Transform」的拓扑里实现,而不是各自拉一份数据。合并后同一份数据只消费一次,下游各分支按自己的逻辑处理,既能降低 Kafka 压力,也避免了同一份数据被重复计算。
5.2 数据倾斜:单 Subtask 积压,整体延迟升高
现象:作业整体并行度正常,但某个 Subtask 的 processing rate 持续偏低,导致 Kafka 消费位点被拉远,数据延迟从分钟级涨到小时级。
原因:实时数仓里常见的倾斜来源是 join 的 key 分配不均——比如按用户 ID join,头部用户的数据量远大于长尾用户;另一个来源是 Kafka 消息 key 的 hash 不均,导致某个分区的数据量特别大。Flink SQL 里 group by 的 key 分布不均就会造成「热点 Subtask」。
解决:第一层用加盐(salting)打散热点 key。在 group by 前对 key 拼接随机后缀,先在打散后的粒度做局部聚合,再去掉后缀做全局聚合:
-- 有倾斜的写法:直接按 user_id 聚合 SELECT user_id, COUNT(*) FROM click_events GROUP BY user_id -- 打散方案:先对 user_id 加盐做局部聚合,再按真实 user_id 汇总 SELECT user_id, SUM(cnt) FROM ( SELECT user_id, COUNT(*) AS cnt FROM ( SELECT user_id, CONCAT(CAST(user_id AS STRING), '_', FLOOR(RAND() * 10)) AS salted_key FROM click_events ) GROUP BY user_id, salted_key )salt_agg GROUP BY user_id第二层调大并行度,同时要保证 Kafka topic 分区数足够多,否则并行度上去了但 consumer 拿不到更多分区。第三层如果倾斜源来自维表 join 的 hot key,把维表数据预加载进内存并使用异步 IO,可以缓解维表查询热点。
5.3 Checkpoint 一直失败或超时
现象:作业运行一段时间后,checkpoint 频繁失败,日志报Checkpoint expired before completing,状态不增长但作业恢复越来越难。
原因:三类常见原因。一是集群 IO 压力大,checkpoint 写入 HDFS 的耗时增长,超过超时阈值;二是 RocksDB 的状态在 checkpoint 时要做快照,状态大、磁盘慢就会卡住;三是作业里有不支持 checkpoint 的算子,比如某些外部 IO 没有实现快照接口。
解决:优先调整 checkpoint 超时时间和间隔,判断是偶发还是持续。如果持续超时,查 HDFS 写入速率,同时查看 Flink UI 里每个 Subtask 的 checkpoint 耗时分布,定位是哪个算子拖后腿。如果是反压导致的 checkpoint 和数据处理互相争抢资源,优先解决反压,而不是盲目加大 checkpoint 频率。我常用的组合:RocksDB 状态后端 + 增量 checkpoint,超时设 5 分钟,间隔 60 秒。
5.4 作业重启后的状态恢复与数据重复
现象:作业因代码逻辑 bug 重启,重启后数据出现重复或丢数据,而且没有统一的去重层兜底。
原因:Flink 默认的 exactly-once 语义在有外部依赖时,只能保证 Flink 内部状态的一致性,不能保证下游 Kafka 或 MySQL 在重复写入时的幂等性。作业从最近一次 checkpoint 恢复,这段窗口内的数据会被重新消费一遍,如果下游没有去重,就会重复。
解决:给下游存储加幂等能力。Kafka sink 开启 idempotent producer(Flink Kafka connector 默认开启),两张方式结合:一是下游消费端按业务主键去重——实时数仓里最常见的是在 DWD 层对 ODS 数据做去重时用ROW_NUMBER() OVER (PARTITION BY pk ORDER BY ts DESC),只保留每个主键的最新一条;二是写入 MySQL/ES 的使用 upsert 语义,而不是 insert。OPPO 在 DWD 层肯定做了这类处理,否则长时间运行下去重复数据会污染所有下游报表。
5.5 业务变更引起的 topic schema 演进
现象:上游业务方在 Kafka topic 里加了字段,下游 Flink SQL 作业直接解析失败,作业重启。
原因:Flink SQL 在建表语句里定义了 schema,topic 里新字段没有对应定义,格式校验失败。
解决:约定一套 schema 演进规则。字段只做 append 不加删除,新增字段默认值代填;Flink SQL 表定义中的字段用ROW包裹便于向后兼容;消息格式用 Avro 而不是 JSON,因为 Avro 自身带 schema 演进机制,JSON 不具备。实时数仓要做到「下游不宕机、字段能自愈」,这需要采集端和上游业务方对 schema 生命周期有强约束。
6. 进阶实践:StreamGraph 重写的具体写法与收益验证
实时数仓的优化的最后一公里,往往不在 SQL 逻辑本身,而在 Flink 生成的执行图上。以「Kafka topic 重复消费」优化为例,OPPO 提到重写 StreamGraph、合并重复 DataSource,这里我拆一个具体的验证流程,你可以直接照做。
在 Flink SQL 作业提交前,用TableEnvironment.getPlanUpdater()(Flink 1.13+)或者通过TableEnvironment.explainSql()先看执行计划。我会先写一个小工具类,把 StreamGraph 中所有 Source 节点打印出来:
// 伪代码:遍历 Flink StreamGraph,找出所有 Kafka source 节点 StreamGraph streamGraph = env.getStreamGraph(true); for (StreamNode node : streamGraph.getStreamNodes()) { if (node.getOperatorFactory() instanceof SourceStreamOperatorFactory) { SourceStreamOperatorFactory<?> factory = (SourceStreamOperatorFactory<?>) node.getOperatorFactory(); // 从 factory 中拿到 Kafka topic、group id 信息,按 topic+group 分组 System.out.println("source node: " + node.getId() + ", parallelism: " + node.getParallelism() + ", operator: " + factory.getClass().getSimpleName()); } }这段代码帮你建立「看到执行图」的意识。生产里很多 Flink SQL 作业的隐性性能问题,比如重复 source、重复 filter、无法并行的维表 join,都能通过explainSql()提前看出来。把执行计划中逻辑节点数量和实际 SQL 条数对比,如果逻辑节点数明显多于预期,审视一下有没有可以合并的算子。
做 DataSource 合并时,注意一个前提:只有 schema 完全相同、消费配置完全相同(group id、offset 模式、topic)的 source 才能安全合并。否则强行合并会让不同语义的 SQL 互相污染状态。我在实际项目里,一般只合并同一份原始日志 topic 的多个读,不合并不同业务 topic。
从那以后,我每次提交 Flink SQL 作业,都强制走一遍四步:第一步explainSql()看执行计划;第二步检查 StreamGraph 里有没有重复 source;第三步看 Checkpoint 间隔和并行度是不是匹配;第四步模拟上游 Kafka 停写 5 分钟测重启恢复。这套流程帮我拦下了至少七成线上事故。希望帮到你。
OPPO 这份实践最值得吸收的,不是某个具体参数,而是「把实时数仓当作系统工程来建设」的态度:统一采集、统一元数据、统一 SQL、统一调度,最后才谈得上稳定的实时服务。你落地时不用一步到位,先从「统一 SQL + 分层建模」开始,跑通后再逐步补管理能力。
本文还有配套的精品资源,点击获取