做数据开发这些年,我身边几乎每隔一段时间就会听到类似的问题:为什么后台报表数据还停留在昨天?为什么运营要的实时销售排行一直出不来?为什么同一个订单,在订单系统里已经付完款,数仓里却还是“待支付”?
这些问题的根源,其实就一句话:离线数仓的 T+1 节奏已经满足不了业务对时效性的要求。于是“实时数仓”就从概念走向了工程落地。而实时数仓里最容易让人头疼的一环,就是把散落在多张业务表里的数据同步成一张可以直接用的宽表——也就是所谓的“宽表同步”。
这篇文章,我打算把自己在实时数仓项目里的实际经验做一个完整复盘。不会只讲概念,也不会只贴代码。会从架构思路、技术选型、宽表同步的几种常见方案,到手把手搭一条订单宽表的实时链路,再到上线后那些逃不掉的坑,一次性讲透。适合正在做实时数仓开发、或者准备把核心报表从离线切到实时的朋友参考。
1. 先把实时数仓的架构思路理顺
1.1 实时不是把离线任务跑快点
很多刚接触实时数仓的同学会有一种错觉:实时数仓就是给离线任务换一个更快的引擎,把每天跑一次的调度改成每小时甚至每分钟跑一次。这个理解不算全错,但会漏掉实时数仓最核心的东西。
离线数仓处理的是“某个时间点已经稳定的数据”,你跑的是批量计算,数据本身就是静止的,算错了重跑一批就行。实时数仓处理的是“还在持续发生的事件流”,数据是被一条一条消费进来的,且天然存在乱序、重复、延迟到达这些离线场景几乎不用考虑的问题。再加一句扎心的:实时任务一旦有问题,你不能像离线一样“今天跑了错,明天重跑”就结束,错误的变更可能已经写到下游、被业务看到、甚至触发了后续动作。
所以实时数仓不是“跑得更快的离线数仓”,它是一套需要同时解决事件捕获、流式计算、状态管理、幂等写入、数据回溯的完整工程体系。这里的核心矛盾在于:源数据在 MySQL 等 OLTP 系统里,目标查询在 OLAP 系统里,中间还隔着加工环节。你要做的就是把这套链路稳定地跑起来,而不是简单地“把速度提上去”。
我见过不少项目,一上来就急着买引擎、搭 Flink,结果跑了两周发现数据对不上,查来查去是源头 binlog 里的字段没解析全。原因就是只关注了“快不快”,没关注“稳不稳、对不对”。做实时数仓,顺序应该是先把数据链路理清楚,再谈性能优化。
1.2 实时数仓的分层和离线有什么不一样
离线数仓的分层大家都很熟:ODS、DWD、DWS、ADS。实时数仓沿用这套分层思想,但每一层的实现方式差别很大。
ODS 层在实时链路里通常就是 Kafka 里的 Topic。你在 MySQL 上开启 binlog,通过 Flink CDC 把变更事件写到 Kafka,这就构成了贴源层。注意 ODS 层不要做太多加工,尽量保留 binlog 原始语义,也就是记录每行数据变更前和变更后的值,这样后续链路才有重试和回溯的空间。
DWD 层是实时数仓的主战场,也是宽表最常出现的地方。这一层要做的事情是清洗、标准化字段、补全维度信息,把多张表的数据合并成一张可以直接查的明细宽表。离线里 DWD 层以 Hive 表为主,实时链路里则往往落在一张 Doris 或 StarRocks 等 OLAP 表上。
DWS 层做轻度汇总,服务大屏、看板这类场景,比如按小时粒度统计销售额、按城市统计订单量。ADS 层更贴近具体应用,一个报表一张表,甚至一个接口一张表。
这个分层里有一个关键差异:离线数仓因为数据是分批落库的,天然有明确的时间边界;实时数仓的数据是持续流动的,所以每一层都必须支持停机续跑和数据回溯。这也是为什么很多实时任务会把中间结果写进 Kafka,而不是只在最终目标表里留一份。Kafka 在这里扮演的角色,相当于离线数仓里那个“可重跑的分区”。
1.3 引擎到底怎么选:以 Doris 为中心的分析
实时数仓的引擎选型,直接影响后面半年你写作业的舒服程度。我评估过 ClickHouse、Doris、StarRocks,也短暂碰过 HBase 和 ES 方案,最终稳定下来的组合是 Kafka + Flink + Doris。下面把这几类引擎的适用边界说一下。
ClickHouse 的宽表聚合分析能力确实顶,列式压缩也做得很好,但它在多表关联场景下比较吃力,高并发点查也不是它的强项。如果你的业务是“一张超大宽表做聚合报表”,ClickHouse 很合适;如果宽表本身要承接多张源头表的变更合并、还要支持按主键高频更新,ClickHouse 会写得很别扭。
HBase 和 ES 就不展开细说了,一个擅长 KV 写入但 SQL 分析能力弱,一个擅长全文检索但聚合分析一般。宽表同步这个场景里,它们要么做存储底座时需要再套一层查询引擎,要么根本不适合作为实时数仓的核心存储。
我最后选 Doris,核心原因有三个。第一,Doris 的主键模型(Unique 模型)天然支持高并发 upsert,这对宽表同步非常重要——你从多张源表拿到的变更最终都要按主键落进同一张表,Doris 在存储层直接处理 merge,省掉了计算层的复杂 join。第二,Doris 对 Flink 的写入生态比较完善,官方 connector 封装了 stream load,配合 checkpoint 机制能做到可靠的幂等写入。第三,Doris 的运维负担相对小,FE + BE 两个角色,集群规模可控,不需要额外引入太多组件。
选型这块我的建议是别迷信某个引擎的“跑分”,把你最核心的两个场景列出来:一是宽表写入的并发和延迟,二是下游查询的复杂度和 QPS。拿真实数据量去压测,比看任何技术测评都靠谱。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 宽表同步的几条常见技术路线
2.1 双流 Join:最直观,但状态开销要算清楚
宽表同步本质上要解决一个问题:多张表的数据怎么合成一张表。最直觉的思路是,在 Flink 里把两条流 join 在一起。比如订单主表和订单明细表,用订单 ID 做关联,输出结果写进目标宽表。
这种做法的优点是链路短、语义直观。但代价是 Flink 的 regular join 需要在状态里缓存已经到达但还没有被匹配上的数据。比如订单主表先到了,明细表的数据晚 10 分钟才到,那这 10 分钟内主表数据必须在状态里待着。如果两条流的延迟波动大,状态就会越积越大,最后要么内存扛不住,要么引入 RocksDB 后吞吐下降。
Flink 1.15 之后 regular join 的实时语义才比较完善,在此之前很多人用 interval join 或者自己维护状态。interval join 限定两条流的数据必须在一个时间窗口内匹配,能控制状态大小,但前提是业务场景本身允许窗口限制。订单主表和明细表这种关联,两边数据到达时间通常相差几秒,窗口可以开得比较宽;但如果是订单和物流状态这种可能隔好几个小时才更新的关联,窗口就不好设了。
我踩过的坑是:为了控制状态给 regular join 设了太短的 TTL,结果部分匹配不到的数据直接过期,宽表里出现了空字段,业务对账才发现丢了数据。状态 TTL 不是越大越好,也不是越小越省,要根据你最慢的那一路数据的延迟来定。
2.2 Lookup Join:维度补全的常规武器
宽表里除了业务事实数据,还有大量的维度信息,比如店铺名称、用户昵称、商品分类。这些维度数据的特点是可以从维表里实时查,不需要流式 join。Flink 的 lookup join 就是干这个的。
lookup join 的写法很简单,在 SQL 里用 FOR SYSTEM_TIME AS OF 去关联一张 JDBC 维表。它的执行方式是:主表的一条数据到了,去查一次维表,把维表字段补上再输出。所以性能瓶颈在于查维表的频率和延迟。
如果你直接用 JDBC connector 不配置任何缓存,每一行数据都会打一次数据库,QPS 一上来维表数据库基本就扛不住了。所以我通常会把维表缓存打开,比如设置 lookup.cache.max-rows 和 lookup.cache.ttl,让一批行共享缓存查询结果。代价是维表数据更新后,缓存里的旧值还要再存活一段时间,宽表里的维度字段可能有短暂延迟。这个延迟你能接受就缓存开大点,不能接受就把 TTL 调小,或者干脆不用 lookup join。
还有一个更稳的思路:如果维表本身数据量不大、变化也不频繁,直接用 Flink CDC 把维表同步到一个本地状态或内存结构里,查询就变成了本地查 Map,比任何外部 lookup 都快,也不存在缓存过期的问题。这条思路在实时数仓开发里很实用,尤其是店铺、商品这类相对稳定的维度。
2.3 主键表 Upsert:我更推荐的多源合并方案
双流 join 和 lookup join 各有适用场景,但如果你想把多张业务表合成一张宽表,还有一条更稳的路线:主键表 upsert。
思路不复杂。每张源头表的数据,各自经过 CDC 进 Kafka,再各自有一个 Flink 任务去消费、清洗,然后写入目标宽表。目标宽表使用主键模型,Doris 在存储层根据主键自动做 merge。也就是说,订单主表的数据到了会更新这一行的订单状态字段,用户维表的数据到了会更新这一行的用户名字段,两者互不干扰。
这个方案最大的优势是解耦。每一张源表的数据链路是独立的,不再需要把所有数据汇集到一个 Flink job 里做 join,状态管理简单得多,故障爆炸半径也小。某张源表的 CDC 挂了,只需要修复这一条链路,其他表还能继续写。
但这里有一个前提必须说清楚:主键表 upsert 要求宽表的粒度是唯一的。如果宽表一行代表一个订单,那订单明细怎么办?你需要先做设计上的选择:要么宽表粒度做到订单明细级别,每个明细一行,订单主表字段冗余在每一行里;要么在源头先做聚合,比如把订单的总金额、总件数算好,再以订单号为粒度同步。
我们线上实践下来,绝大多数宽表都采用了“明细级别一行”的方式,因为这样既保留了明细查询能力,又能在 DWS 层再按需聚合。主键就用明细表的自增 ID,订单号作为普通索引列。这样设计的好处是:任何一张源头表的数据更新,都只影响宽表的某一行,Doris 的 merge 压力小,也不会出现两个任务并发写同一行导致字段互相覆盖的情况。
2.4 我的推荐组合
我最终沉淀下来的宽表同步方案是这样的:所有业务表通过 Flink CDC 采集,统一入 Kafka,按表名建 Topic,这是 ODS 层。然后一个“清洗 + 维度退化”的 Flink 任务,消费 ODS 数据,做字段标准化、格式转换、必要时用 lookup join 补维度,输出成宽表口径的变更流,写入另一个 Kafka Topic。最后再有一个纯 sink 任务,消费这个宽表 Topic,写入 Doris。
多出来的这层 Kafka,就是整个链路的安全垫。Flink 任务重启时不会导致目标表写一半丢一半,数据回溯时也不用重建整条链路,直接从 Kafka 重放即可。延迟会多那么一两秒,但相比稳定性收益,这点延迟完全可以接受。
3. 订单宽表:从 CDC 到 Kafka 再到 Doris 的完整实操
3.1 场景定义和宽表结构设计
纸上谈兵没意思,下面我用一个真实的订单宽表场景,把全链路参数和 SQL 过一遍。
假设业务库有三张核心表:order_main(订单主表,包含订单号、用户 ID、店铺 ID、订单状态、支付金额、下单时间、更新时间);order_item(订单明细,一个订单多件商品,包含明细 ID、订单 ID、商品 ID、单价、数量);还有两张维度表 shop(店铺)和 user(用户)。
需求方要的宽表粒度是订单明细级别,每一行是一个商品,同时包含订单属性、店铺属性、用户属性。这张宽表要用在实时订单列表、大促期间的商品销售排行、客服的订单轨迹查询。
DDL 设计如下:
sql复制CREATE TABLE dwd_order_wide (
id BIGINT NOT NULL,
order_id BIGINT NOT NULL,
order_no VARCHAR(64),
user_id BIGINT,
user_name VARCHAR(64),
shop_id BIGINT,
shop_name VARCHAR(128),
goods_id BIGINT,
goods_name VARCHAR(256),
price DECIMAL(12,2),
quantity INT,
gmv DECIMAL(12,2),
order_status TINYINT,
create_time DATETIME,
update_time DATETIME
) ENGINE=OLAP
UNIQUE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 24
PROPERTIES (
"replication_num" = "3",
"enable_unique_key_merge_on_write" = "true",
"storage_format" = "v2"
);
几个细节解释一下。主键用明细 ID 而不是订单号,原因是明细 ID 才符合“一行一个商品”的粒度,而且数字主键做 hash 分布更均匀,不会让某个分桶的数据量畸高。分桶数 24 是我按预估千万级数据量、每桶几百万行的水平定的。分桶太少会热点,太多会造成小文件问题,经验值是在数据量增长后再扩容,而不是一开始就贪多。
enable_unique_key_merge_on_write 必须打开,这样 Doris 在写入时就能实时合并同主键数据,否则要通过读时合并,对高频 upsert 场景不友好。
3.2 建表与通道配置:Doris DDL 和 Kafka Topic
Kafka 侧,我建了三个 Topic:ods_order_main、ods_order_item 各存源表变更事件,dwd_order_wide_upsert 存宽表变更流。Topic 分区数设为 24,和 Doris 分桶数对齐,这样 Flink 并行写 Doris 时能尽量分散到不同分桶,减少写入热点。副本数设 2,数据重要程度高且对延迟要求苛刻的 Topic 可以设 3。
这里有一个很多人忽略的点:Flink CDC 直接写 Kafka 时,使用 upsert-kafka connector,必须把主键字段同时作为消息 Key。这样同一行数据的变更事件都会进入同一个 Kafka 分区,保证该行的变更顺序不乱。如果 Key 设置错了,同一主键的变更被 hash 到不同分区,Flink 消费后可能出现先更新后插入的乱序,下游 Doris merge 之后的数据就是错的。
Doris 侧的连接可以使用 MySQL 协议或 stream load。Flink 官方 Doris connector 默认走 stream load,写起来也比较简单,只要把 fenodes、table.identifier、账号密码、字段映射配好就行。
3.3 Flink SQL 组装宽表,几个参数必须调
源表 CDC 的建表语句,我以订单主表为例:
sql复制CREATE TABLE order_main_cdc (
id BIGINT,
order_no STRING,
user_id BIGINT,
shop_id BIGINT,
order_status TINYINT,
gmv DECIMAL(12,2),
create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'xxx',
'port' = '3306',
'username' = 'cdc_user',
'password' = 'xxxx',
'database-name' = 'trade_db',
'table-name' = 'order_main',
'scan.startup.mode' = 'initial',
'server-time-zone' = 'Asia/Shanghai'
);
scan.startup.mode 这里我解释一下。首次上线建议用 initial,它会把表里已有的历史数据先做一次快照,然后无缝切换到增量 binlog。如果只想要上线之后的增量数据,可以改成 latest-offset。我建议首次初始化时老老实实用 initial,除非这张表数据量已经大到快照会影响源库性能。
server-time-zone 必须配成 Asia/Shanghai,否则时间字段解析出来会差 8 小时,这个问题在时间字段对账时特别隐蔽。
维表用 JDBC connector 加缓存:
sql复制CREATE TABLE dim_shop (
shop_id BIGINT,
shop_name STRING,
PRIMARY KEY (shop_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://xxx:3306/dim_db',
'username' = 'xxx',
'password' = 'xxx',
'table-name' = 'shop',
'lookup.cache' = 'PARTITION',
'lookup.cache.max-rows' = '20000',
'lookup.cache.ttl' = '30min',
'lookup.max-retries' = '3'
);
这里 lookup.cache 设为 PARTITION,意思是每个并行子任务维护一份缓存,避免并发访问同一个外部数据库连接。TTL 设 30 分钟,店铺改名、上下架状态这类维表更新不频繁,30 分钟的缓存延迟业务可接受。如果你做的是价格、库存这类高实时维度,TTL 就要缩到几秒甚至不用缓存。
最终宽表写入 Doris:
sql复制CREATE TABLE dwd_order_wide_sink (
id BIGINT,
order_id BIGINT,
order_no STRING,
user_id BIGINT,
user_name STRING,
shop_id BIGINT,
shop_name STRING,
goods_id BIGINT,
goods_name STRING,
price DECIMAL(12,2),
quantity INT,
gmv DECIMAL(12,2),
order_status TINYINT,
create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'doris',
'fenodes' = '127.0.0.1:8030',
'table.identifier' = 'dwd.dwd_order_wide',
'username' = 'xxx',
'password' = 'xxx',
'sink.label-prefix' = 'dwd_order_wide',
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true',
'sink.batch.size' = '10000',
'sink.batch.interval' = '5s'
);
组装宽表的 INSERT 语句:
sql复制INSERT INTO dwd_order_wide_sink
SELECT
i.id,
m.id AS order_id,
m.order_no,
m.user_id,
u.user_name,
m.shop_id,
s.shop_name,
i.goods_id,
g.goods_name,
i.price,
i.quantity,
m.gmv,
m.order_status,
m.create_time,
m.update_time
FROM order_item_cdc i
LEFT JOIN order_main_cdc m
ON i.order_id = m.id
LEFT JOIN dim_shop s
FOR SYSTEM_TIME AS OF i.proc_time
ON m.shop_id = s.shop_id
LEFT JOIN dim_user u
FOR SYSTEM_TIME AS OF i.proc_time
ON m.user_id = u.user_id;
这个链路里,订单明细是主表,明细来了必须立刻查出对应的订单信息。这里 order_main_cdc 依然用的是双流 join 方式,但因为它和明细变更的时间差通常在秒级,我把 Flink 的 table.exec.state.ttl 设成了 2 小时,既能保证明细到的时候订单数据还在状态里,也不会撑爆 RocksDB。
Flink 里我建议把 checkpoint 间隔设为 60 秒,连续失败容忍次数设为 3 次。实时任务的核心目标不是“永远不挂”,而是“挂了能快速恢复且不丢数据”。checkpoint 太频繁会影响吞吐,太稀疏会增加恢复时间,60 秒是目前比较稳妥的平衡点。
3.4 上线前的数据校验与回放
上线实时链路前,最容易被忽略但最不能省的是数据校验。我在项目里总结了一套三层校验法。
第一层是数量校验。取源表 order_item 的 count,对比 Doris 宽表里同一时间段的 count。实时环境里两边永远不可能完全一致,因为时刻在写入,但差值应该稳定在一个很小的波动范围。
第二层是字段校验。在源表里随机抽 100 个订单 ID,把对应的订单主表、明细、维度表数据都查出来,跟 Doris 宽表里同订单的数据逐字段比对。重点看时间字段有没有偏时区、金额有没有精度丢失、状态映射是否正确。
第三层是链路校验。用 Kafka 客户端直接消费 dwd_order_wide_upsert 这个 Topic,看变更消息是否完整,有没有大量重复或乱序。这一层是刻意做的,因为 Doris 里最终结果是对的,不代表中间流是健康的。中间流一旦有问题,Doris 数据迟早会出问题,早发现早处理。
校验通过后,再切线上流量。还有一点要特别提醒:如果多个 CDC 源表需要构造同一张宽表,要先确认它们快照的起始时间一致,否则可能出现订单明细读到了、订单主表的快照还没有读到,导致明细宽表的订单字段为空。我自己遇到过这种情况,排查半天,最后发现是两个 CDC job 启动时间差了几分钟,快照边界错位导致。解决办法是让这些源表任务从同一个时间点开始消费,或者对订单主表这类关键表做一次补数据。
4. 上线后那些逃不掉的坑
4.1 时间不对、数据对不上,先查源头
线上跑一段时间后,业务方反馈数据不对是最常见的。我列一个排查顺序,基本能覆盖 80% 的问题。
先看时间字段。这是最常见也最坑的。Flink CDC 读 MySQL 时如果不配 server-time-zone,默认按服务器时区解析,经常出现整张宽表时间偏 8 小时的问题。表面上看每条数据都在更新,但时间就是不对。直接查源表和目标表同一条订单的 create_time,一眼就能看出来。
再看源表快照。如果任务启动时用了 latest-offset,历史数据没进来,宽表里就会一直缺旧数据,但增量数据又是正常的。这种问题需要在需求阶段就跟业务对齐:是要全量历史加增量,还是只要增量。
然后是主键重复。如果宽表主键设计得不合理,比如用业务订单号而不是明细 ID,同一个订单多个明细就会互相覆盖,最终只剩一条。这种问题的特征是:count 明显小于源表明细数,订单号出现次数永远为 1。这也是我反复强调宽表粒度设计必须前置、不能上线后改结构的原因。
4.2 状态爆炸和延迟,根因不一定在 SQL
实时任务跑一段时间后,Flink Web UI 里的状态存储越来越大,最终导致 GC 频繁,延迟飙升。这种情况我遇到的第一反应不是去动 SQL,而是先分清状态到底在哪个算子。
Flink Web UI 的 Task Manager 页面看每个算子的 state 大小。如果是 join 算子的状态大,说明你的双流 join 或 interval join 里,某路数据到达时间差太大,或者缓存了太多没有匹配上的行。这时候调 TTL 是最快的办法,但要注意不能一味缩小 TTL 导致丢数据。比较稳妥的做法是复盘一下业务上这两路数据是否有明确的时间关联,如果有,改用 interval join 把窗口收紧;如果没有强关联,考虑换成主键 upsert 方案,在存储层合并,而不是在计算层 join。
如果是窗口算子的状态大,比如做 1 小时的滚动聚合,数据量又大,可以考虑加细粒度预聚合,或者换用 Flink SQL 的 emit 策略。另外,确认一下是否把状态后端换成了 RocksDB,堆内状态在大数据场景下很难撑住。
4.3 写入抖动的排查思路
Doris 写入抖动最常见的原因是 stream load 的 label 冲突或批次配置不合理。
Flink Doris connector 默认用 sink.label-prefix 加时间戳做 stream load 的 label。任务重启后从 checkpoint 恢复时会沿用之前的 label,保证写入幂等。但如果你的 sink.label-prefix 在多个任务里重复了,Doris 会拒绝一部分写入,表现就是任务日志里出现 label 相关的异常。排查方法很简单:确认每个 DWD 任务的 label-prefix 全局唯一。
批次参数也要注意。sink.batch.size 和 sink.batch.interval 要配合数据吞吐来设。数据量小的时候,批大小设太大反而会导致数据在内存里攒太久,延迟变高。数据量大的时候,批大小设太小会产生大量高频 stream load,Doris 端小文件变多,查询性能下降。我一般把批大小设在 1 万行、批间隔 5 秒,再根据线上写入速率微调。
如果你的 Doris 版本支持 group commit 特性,可以开启,让写入侧更平滑,减少小文件。这个特性对不同版本支持情况不一样,使用时先确认版本兼容性。
4.4 实时任务的数据修复套路
实时数仓最头疼的就是上线后才发现业务口径定义错了,或者之前有一段时间的数据写坏了。离线任务直接回溯重跑历史分区就行,实时任务没有“历史分区”,数据都是流式覆盖的,修复起来要讲究方法。
我常用的套路是 Kafka 重放。前提是你在链路里保留了宽表变更流 Topic。修复时,先停掉下游 sink 任务(或者只停问题数据涉及的消费组),定位到错误时间段的 Kafka offset,把这段消息重新消费一遍。如果只是字段计算错误,可以在重放时加一个临时处理逻辑,比如修正某个金额字段再写回。如果错误范围太大,也可以直接从 ODS 层 Kafka 重放,走一遍正常的宽表加工逻辑。
这个套路能成立的前提是:Kafka Topic 的 retention 时间够长。我一般给 ODS 和 DWD 层 Topic 的 retention 设 3 到 7 天,代价是多占存储,但换来的是修复问题的从容。相比数据错了被业务追着问、又只能临时写一个补数任务手动 update Doris,这步投入非常值。
4.5 运维细节:让实时任务少折腾
最后分享几条我常用的运维经验。
第一,Kafka Lag 是实时链路健康度的第一指标。给每个消费组都配上 Lag 监控,Lag 持续增长说明消费能力跟不上,爆炸只是时间问题。第二,Flink 的 checkpoint 时长也是核心监控项,单次 checkpoint 超过 5 分钟基本说明有大对象或者反压,要尽快处理。第三,Doris 的 BE 节点磁盘水位、FE 的内存变化这类元数据健康度,提前做告警,别等查询超时了才被动发现。第四,重要任务在更新 SQL 前,先把旧作业的状态快照做备份,虽然 Flink 本身有恢复机制,但多一手准备总没错。
实话说,实时数仓项目能不能稳,一半靠架构,一半靠运维习惯。架构上多用 Kafka 做缓冲、少在计算层维护大状态、主键表 upsert 优先,这些决定了系统的上限;运维上监控到位、报警及时、预案充分,决定了你凌晨三点会不会被电话叫醒。这两半都补齐了,实时数仓才能真正成为可信赖的数据底座,而不是那个“看起来很快但总让业务不放心”的临时方案。
