简介:这份PDF文档面向数据架构师、实时计算开发者和数仓工程师,系统讲解如何基于Flink与Hologres构建云原生实时数仓。内容围绕Lambda架构的局限展开,深入剖析HTAP与HSAP理念,涵盖实时导入与批量归档、维表关联与离线加速、联邦计算、结果缓存等关键实践,并详解Hologres计算存储分离、流批统一存储、C++ Native执行引擎与优化器等底层设计,帮助读者理解实时离线一体化与分析服务一体化的落地路径。资源包共1个PDF文件,大小约1.23MB,轻量便携,适合通勤或碎片时间研读。目前已有593人学习下载,文档以架构图与要点式讲解为主,逻辑清晰,可作为技术选型与方案设计的参考手册,帮助读者快速把握Flink与Hologres组合在实时数仓场景中的核心思路与优化方向。
1. Flink + Hologres 云原生实时数仓:从 CDC 入仓到 OLAP 查询的完整链路
电商大促的凌晨,订单库的 Binlog 以每秒几万条的速率往外涌,下游既要实时看 GMV 大盘,又要支撑运营按 SKU、按渠道、按省份做多维下钻。传统做法是 T+1 把数据抽到离线数仓,第二天早上才能出报表,运营等不了;换成 Flink 消费 Kafka 再写进 Hologres,链路能压到秒级,但真正落地时会撞上一堆问题:CDC 全量阶段把源库拖垮、维表关联把内存打爆、写入 Hologres 时连接数不够、Exactly-Once 配了却还是重复。这篇笔记就围绕 Flink Hologres 云原生实时数仓这条链路,把选型理由、建表 DDL、CDC 配置、写入参数和排查手段讲清楚,适合正在做实时数仓选型、或者已经上了 Flink 但写入侧不稳的工程师。
2. 为什么是 Flink 加 Hologres:云原生实时数仓的选型账
2.1 实时数仓的三种技术路线对比
在动手之前,先把路线选清楚。实时数仓的存储层大致有三条路:一是 Flink 直接写 HBase 或 Redis,查询能力弱,只适合点查;二是 Flink 写 ClickHouse,写入吞吐高、单表查询快,但多表 Join 和更新场景吃力,CDC 的 Upsert 语义要靠 ReplacingMergeTree 绕;三是 Flink 写 Hologres,走 Binlog 订阅式的行存加列存混合,天然支持 Upsert 和主键更新,还能直接对接 MaxCompute 做湖仓一体。
| 维度 | HBase/Redis | ClickHouse | Hologres |
|---|---|---|---|
| 写入语义 | Put,无更新语义 | 追加为主,更新靠合并 | 主键 Upsert,支持部分列更新 |
| 多表 Join | 不支持 | 大表 Join 受限 | 支持,可下推 |
| 点查延迟 | 毫秒级 | 毫秒到秒级 | 毫秒级 |
| 与离线打通 | 弱 | 需导出 | 直读 MaxCompute |
| 运维成本 | 自建集群 | 自建或云托管 | 云原生托管 |
选 Hologres 的核心理由是它把「实时写入」和「分析查询」放在同一个引擎里,不需要再维护一条从实时到离线的同步链路。云原生在这里的价值不是概念,而是存储计算分离之后,写入节点和查询节点可以独立扩缩,大促前把计算组拉起来,结束后缩回去,成本可控。
2.2 云原生架构下 Flink 与 Hologres 的分工边界
分工要划清楚,否则后面调优会互相甩锅。Flink 负责的是「流式加工」:CDC 解析、维表关联、窗口聚合、脏数据分流。Hologres 负责的是「存储与服务」:主键去重、列存压缩、索引加速、对外提供 JDBC 查询。两者之间通过 Hologres 的 Flink Connector 通信,写入走的是 Hologres 的实时写入接口,不是 JDBC 批量 Insert。
一个常见的误区是把聚合逻辑全压在 Hologres 侧,用物化视图或定时刷新来做。实时场景下,聚合应该在 Flink 的窗口里完成,Hologres 只存结果宽表。原因很简单:Flink 的状态后端可以扛住乱序和迟到数据,Hologres 的查询资源要留给下游 BI 和 Ad-hoc 查询,不该被写入侧的聚合拖累。
提示:如果业务方要求「任意维度实时下钻」,不要试图用一张宽表满足所有维度,正确做法是 Flink 侧拆成多个轻度聚合的 DWS 表,Hologres 侧用 Join 组合。
3. 从 MySQL Binlog 到 Hologres 宽表:CDC 入仓的最小可跑通链路
3.1 环境准备与依赖版本对齐
先确认版本。Flink 用 1.17 或 1.18 比较稳,Connector 版本必须和 Flink 大版本对齐,否则会出现类加载冲突。Hologres 的 Flink Connector 在 Maven 中央仓库可以拉到,CDC 用 flink-connector-mysql-cdc。下面是一个最小 pom 依赖片段。
<!-- Flink 1.17 + Hologres Connector + MySQL CDC --> <dependency> <groupId>com.alibaba.hologres</groupId> <artifactId>hologres-connector-flink-1.17</artifactId> <version>1.6.0</version> </dependency> <dependency> <groupId>com.ververica</groupId> <artifactId>flink-connector-mysql-cdc</artifactId> <version>2.4.2</version> </dependency>版本对齐的逻辑:Hologres Connector 的 1.6.x 系列对应 Flink 1.17,1.7.x 对应 Flink 1.18。CDC 2.4.x 支持 Flink 1.17 的增量快照框架。如果版本错配,典型报错是NoSuchMethodError或ClassNotFoundException,排查时先看flink-sql-connector和flink-connector是否混用。
3.2 Hologres 侧建表:主键、分布键与索引怎么定
建表是整条链路的地基。Hologres 建表要关注三件事:主键、分布键(Distribution Key)、聚簇索引(Clustering Key)。主键决定 Upsert 语义,分布键决定数据落在哪个 Shard,聚簇索引决定范围查询的裁剪效率。
-- Hologres 侧订单宽表 BEGIN; CREATE TABLE public.dws_order_wide ( order_id BIGINT NOT NULL, user_id BIGINT, sku_id BIGINT, channel TEXT, province TEXT, pay_amount NUMERIC(18,2), order_status TEXT, update_time TIMESTAMPTZ, PRIMARY KEY (order_id) ); CALL set_table_property('public.dws_order_wide', 'distribution_key', 'order_id'); CALL set_table_property('public.dws_order_wide', 'clustering_key', 'update_time'); CALL set_table_property('public.dws_order_wide', 'segment_key', 'update_time'); CALL set_table_property('public.dws_order_wide', 'bitmap_columns', 'channel,province,order_status'); CALL set_table_property('public.dws_order_wide', 'dictionary_encoding_columns', 'channel,province,order_status'); COMMIT;参数说明:distribution_key选 order_id,保证同一订单的更新落到同一 Shard,避免跨 Shard 更新;clustering_key选 update_time,让时间范围查询能裁剪文件;bitmap_columns给低基数列建位图索引,channel、province 这类枚举值查询会快很多;dictionary_encoding_columns做字典编码,压缩存储。注意主键列不能做 dictionary encoding,会报错。
3.3 Flink SQL 作业:CDC Source 与 Hologres Sink 的完整写法
下面是一个可以直接提交到 Flink 集群的 SQL 作业,从 MySQL 订单表读 CDC,做简单清洗后写入 Hologres。
-- Source: MySQL CDC CREATE TABLE mysql_order ( order_id BIGINT, user_id BIGINT, sku_id BIGINT, channel STRING, province STRING, pay_amount DECIMAL(18,2), order_status STRING, update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'mysql-host', 'port' = '3306', 'username' = 'cdc_user', 'password' = '******', 'database-name' = 'order_db', 'table-name' = 't_order', 'server-time-zone' = 'Asia/Shanghai', 'scan.incremental.snapshot.enabled' = 'true', 'scan.incremental.snapshot.chunk.size' = '8096', 'debezium.snapshot.locking.mode' = 'none' ); -- Sink: Hologres CREATE TABLE hologres_sink ( order_id BIGINT, user_id BIGINT, sku_id BIGINT, channel STRING, province STRING, pay_amount DECIMAL(18,2), order_status STRING, update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'hologres', 'dbname' = 'realtime_dw', 'tablename' = 'dws_order_wide', 'username' = 'access_id', 'password' = 'access_key', 'endpoint' = 'hgprecn-cn-xxx.hologres.aliyuncs.com:80', 'jdbcWriteBatchSize' = '1024', 'jdbcWriteFlushInterval' = '3000', 'connectionSize' = '5', 'mutateType' = 'insertorupdate', 'ignoreDelete' = 'false' ); INSERT INTO hologres_sink SELECT order_id, user_id, sku_id, channel, province, pay_amount, order_status, update_time FROM mysql_order;逻辑说明:CDC Source 用增量快照模式,scan.incremental.snapshot.chunk.size控制全量阶段每批读多少行,8096 是经验值,太小会导致快照阶段慢,太大对源库压力大。debezium.snapshot.locking.mode设为 none,避免全量阶段锁表,这是生产环境必须改的,默认值会加全局锁。
Sink 侧mutateType设为 insertorupdate,对应 Hologres 的 Upsert 语义;ignoreDelete设为 false,保证上游删除能同步下来;jdbcWriteBatchSize和jdbcWriteFlushInterval是一对,前者控制攒批行数,后者控制攒批时间,两个条件谁先满足谁触发写入。connectionSize是写入连接数,一般设 5 到 10,设太大反而会因为连接竞争导致抖动。
注意:Hologres Sink 的
connectionSize不是越大越好。实测在单并发下,5 个连接已经能打满一个 Shard 的写入带宽,加到 20 反而出现连接等待。
4. 写入性能与 Exactly-OOnce:参数调优和状态管理
4.1 攒批参数与 Checkpoint 的配合关系
写入性能的核心矛盾是「攒批大小」和「Checkpoint 间隔」的配合。Flink 的 Checkpoint 会触发 Sink 的 flush,如果 Checkpoint 间隔是 10 秒,而jdbcWriteFlushInterval是 3 秒,那大部分批次是定时触发的,Checkpoint 时只需要 flush 剩余数据,延迟低。反过来,如果 Checkpoint 间隔 1 分钟,攒批时间 30 秒,那每次 Checkpoint 都要等大批次落盘,端到端延迟会飙到分钟级。
推荐配置:Checkpoint 间隔 10 到 30 秒,jdbcWriteFlushInterval设为 Checkpoint 间隔的三分之一到二分之一,jdbcWriteBatchSize设为 1024 到 4096。这样正常流量下靠定时 flush,突发流量下靠攒批行数触发,Checkpoint 时残留数据少。
# flink-conf.yaml 关键项 execution.checkpointing.interval: 15s execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 10min state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: hdfs:///flink/checkpointsstate.backend选 rocksdb 是因为 CDC 的增量快照状态可能很大,内存后端扛不住。state.backend.incremental开启增量 Checkpoint,避免每次全量上传状态。
4.2 维表关联:Lookup Join 的缓存策略与失效
实时宽表通常要关联维表,比如订单关联商品维表拿类目。Flink 的 Lookup Join 会缓存维表数据,缓存策略选错会导致数据不一致或者内存溢出。
CREATE TABLE dim_sku ( sku_id BIGINT, category STRING, brand STRING, PRIMARY KEY (sku_id) NOT ENFORCED ) WITH ( 'connector' = 'hologres', 'dbname' = 'dim_db', 'tablename' = 'dim_sku', 'username' = 'access_id', 'password' = 'access_key', 'endpoint' = 'hgprecn-cn-xxx.hologres.aliyuncs.com:80' ); SELECT o.order_id, o.pay_amount, d.category, d.brand FROM mysql_order AS o LEFT JOIN dim_sku FOR SYSTEM_TIME AS OF o.proc_time AS d ON o.sku_id = d.sku_id;缓存策略通过lookup.cache参数控制,可选 NONE、LRU、ALL。维表小且变更少用 ALL,全量加载到内存;维表大用 LRU,配合lookup.cache.max-rows和lookup.cache.ttl控制。TTL 设太短会频繁查 Hologres,设太长维表变更感知慢。一般 TTL 设 10 分钟,max-rows 设 10 万。
4.3 状态后端与 Checkpoint 的排错要点
状态相关的问题最隐蔽。常见现象是 Checkpoint 一直失败,报Checkpoint expired before completing。原因通常是状态太大,上传 HDFS 超时。解决分三步:先看 Checkpoint 大小,在 Flink UI 的 Checkpoints 页面能看到;如果超过 1GB,检查是否有无界状态,比如 CDC 的增量快照没开、或者窗口没设 TTL;确认状态合理后,调大execution.checkpointing.timeout和state.backend.rocksdb.writebuffer.size。
另一个坑是 RocksDB 的本地目录磁盘满。state.backend.rocksdb.localdir默认在 TaskManager 的临时目录,大状态作业要显式指定到数据盘,并监控磁盘使用率。
5. 避坑与排查:Flink 写 Hologres 最常见的五类翻车
5.1 现象:作业启动后 Hologres 连接数暴涨,报 too many connections
原因:connectionSize设太大,或者作业并发度高,每个并发都建了独立连接池。Hologres 单实例的连接数有上限,默认几百,超了就拒绝。
解决:把connectionSize降到 5 以内,同时用 Hologres 的 Connection Pool 或者把写入并发控制在合理范围。如果并发确实高,考虑在 Sink 前加一层 rebalance 或者用 Hologres 的 Fixed Connection 模式。
5.2 现象:CDC 全量阶段源库 CPU 打满,业务查询变慢
原因:scan.incremental.snapshot.chunk.size设太大,或者debezium.snapshot.locking.mode没改成 none,全量阶段锁表。
解决:chunk size 降到 4096 甚至 2048,加scan.snapshot.fetch.size控制每次 fetch 行数。同时确认 MySQL 的max_connections够用,CDC 会占用源库连接。生产环境建议在从库上做 CDC,不要直接读主库。
5.3 现象:写入 Hologres 报 duplicate key value violates unique constraint
原因:Hologres 表的主键和 Flink Sink 的 PRIMARY KEY 不一致,或者上游数据本身有重复主键但 Flink 没做去重。
解决:核对两边主键定义,必须完全一致。如果上游有重复,在 Flink 侧用ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY update_time DESC)去重后再写入。
5.4 现象:Checkpoint 成功但数据重复写入
原因:mutateType设成了 insert,而不是 insertorupdate。insert 模式下重复数据会直接插入,主键冲突报错或者产生重复行。
解决:确认 Sink 参数mutateType为 insertorupdate,并且 Hologres 表有主键。另外检查 Checkpoint 模式是否为 EXACTLY_ONCE,AT_LEAST_ONCE 模式下重复是预期行为。
5.5 现象:维表关联后数据变多,出现笛卡尔积
原因:维表主键不唯一,或者 Join 条件写错。Lookup Join 要求维表主键唯一,如果维表有重复主键,Flink 会取最后一条,但某些版本会返回多条。
解决:在 Hologres 侧确认维表主键唯一,用SELECT sku_id, COUNT(*) FROM dim_sku GROUP BY sku_id HAVING COUNT(*) > 1排查。Join 条件确保是等值连接,不要用非等值条件。
6. 进阶技巧:用火焰图和 Metrics 定位写入瓶颈
调优到最后,靠猜没用,得看数据。Flink 的火焰图能直接告诉你时间花在哪。在 Flink UI 的 Job 页面点开某个算子,选 Flame Graph,如果发现HologresOutputFormat.flush占比高,说明写入是瓶颈,要调攒批参数;如果Deserialize占比高,说明 CDC 解析慢,要加并发。
Hologres 侧看hg_worker_query_duration和hg_worker_write_rows两个指标,前者是查询耗时,后者是写入行数。如果写入行数远小于 Flink 的输入行数,说明有数据被过滤或者攒批没触发。
我自己的习惯是每次上线新作业,先跑 10 分钟,看三个数:Checkpoint 大小是否稳定、Hologres 写入 QPS 是否匹配输入、端到端延迟是否在预期内。这三个数对了,再放量。有一次大促前没看 Checkpoint 大小,结果状态涨到 8GB,Checkpoint 超时导致作业重启,血泪教训。希望帮到你。
本文还有配套的精品资源,点击获取