ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink + Hologres 云原生实时数仓:CDC 入仓到 OLAP 查询全链路实战

Flink + Hologres 云原生实时数仓:CDC 入仓到 OLAP 查询全链路实战 简介这份PDF文档面向数据架构师、实时计算开发者和数仓工程师系统讲解如何基于Flink与Hologres构建云原生实时数仓。内容围绕Lambda架构的局限展开深入剖析HTAP与HSAP理念涵盖实时导入与批量归档、维表关联与离线加速、联邦计算、结果缓存等关键实践并详解Hologres计算存储分离、流批统一存储、C Native执行引擎与优化器等底层设计帮助读者理解实时离线一体化与分析服务一体化的落地路径。资源包共1个PDF文件大小约1.23MB轻量便携适合通勤或碎片时间研读。目前已有593人学习下载文档以架构图与要点式讲解为主逻辑清晰可作为技术选型与方案设计的参考手册帮助读者快速把握Flink与Hologres组合在实时数仓场景中的核心思路与优化方向。1. Flink Hologres 云原生实时数仓从 CDC 入仓到 OLAP 查询的完整链路电商大促的凌晨订单库的 Binlog 以每秒几万条的速率往外涌下游既要实时看 GMV 大盘又要支撑运营按 SKU、按渠道、按省份做多维下钻。传统做法是 T1 把数据抽到离线数仓第二天早上才能出报表运营等不了换成 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/RedisClickHouseHologres写入语义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 groupIdcom.alibaba.hologres/groupId artifactIdhologres-connector-flink-1.17/artifactId version1.6.0/version /dependency dependency groupIdcom.ververica/groupId artifactIdflink-connector-mysql-cdc/artifactId version2.4.2/version /dependency版本对齐的逻辑Hologres Connector 的 1.6.x 系列对应 Flink 1.171.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_ONCEAT_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 大小结果状态涨到 8GBCheckpoint 超时导致作业重启血泪教训。希望帮到你。本文还有配套的精品资源点击获取
返回列表