Flink+Hologres实时数仓建设实践:从选型到踩坑全记录
做实时数仓这个事很多团队一开始不是被技术难住的而是被“选哪个组件”给难住的。我大概从两年前开始正式碰 Flink 和 Hologres 这条链路当时团队的目标非常朴素报表分钟级延迟DWD 层明细能和离线数仓对齐还要支撑上千 QPS 的在线查询。中间试过不少方案最后沉淀下来的就是 Flink 负责实时计算Hologres 负责明细存储和实时查询。这篇文章把整个项目建设过程完整梳理一遍包括选型思路、表设计、连接器配置、生产 SQL 示例和踩坑记录希望能给正在做实时数仓选型或刚入坑 Flink Hologres 的同学一些参考。1. 实时数仓选型为什么最终是 Flink Hologres1.1 实时数仓到底在解决什么问题离线数仓大家都很熟T1 调度凌晨跑批早上出报表。问题在于业务一旦卷起来T1 就扛不住了。运营要看实时 GMV风控要秒级识别异常行为推荐要基于最近几分钟的点击特征做召回。这些场景本质上是同一件事对持续到达的数据做低延迟的加工然后把加工结果提供给下游查询或决策。实时数仓的难点不在于“实时计算”本身而在于“实时之后怎么查”。Flink 能算但 Flink 不擅长存Kafka 能缓冲但 Kafka 不适合做明细查询和 Ad-hoc 分析传统的 MySQL 存得下数据但大促期间高并发聚合查询会直接把库打挂。所以实时数仓的存储层必须同时具备几个能力写入侧能承接 Flink 的高吞吐流式写入查询侧能扛住高并发点查和复杂 OLAP 查询数据更新上要支持 Upsert 语义否则 Flink 重放数据时会产生大量重复记录。Hologres 正好踩在这些需求点上。它是兼容 PostgreSQL 协议的分布式数仓服务底层是列式存储和 MPP 架构单表支持主键 Upsert也支持行存、列存、行列共存三种存储模式。这意味着同一张表既能做高并发点查又能跑大范围聚合分析不用像某些方案那样为了性能硬拆好几套存储。1.2 和其他组合的对比为什么不是 Kafka ClickHouse也不是 Flink StarRocks选型期我对比过几套主流方案这里说下我自己的判断。Kafka ClickHouse 是很多团队的第一反应。ClickHouse 的聚合查询确实快但它对实时数据更新支持得很别扭主键去重要靠 ReplacingMergeTree重写数据有延迟而且并发点查能力一般。如果你的场景只是“大屏展示 聚合报表”这套组合够用但如果你还要做实时维表关联、明细回溯、ByPrimaryKey 查询ClickHouse 会让人抓狂。Flink StarRocks 也是很成熟的组合StarRocks 在新版本里对主键模型和 Upsert 的支持做得不错。当时没选它主要是因为团队里 PostgreSQL 系的技术栈更熟——Hologres 是 PG 协议BI 工具、SQL 脚本、数据迁移工具基本能无缝复用另外 Hologres 和 Flink 都是阿里系引擎连接器、权限体系、网络打通这几个方面省了很多对接成本。注意这不是说 StarRocks 不好只能说团队技术背景和现有生态决定了哪个更顺手。2. 开工前的三件套链路设计、表模型与权限规划2.1 实时数仓分层建设思路实时数仓不是把离线数仓的表结构搬过来改个名就行它的分层得结合流式处理的特点重新设计。我习惯把链路拆成 ODS、DWD、DWS、ADS 四层但每层的职责和离线有区别。ODS 层承接原始数据整条 Kafka Topic 的数据直接落 Hologres 做暂存或者只做格式校验。这层不复杂但必须保留原始数据方便上游出问题后回溯。DWD 层做清洗、标准化、维表补全。Flink 在这里消费 ODS 层数据做过滤、字段映射、脏数据分流再通过 JOIN Hologres 维表补齐维度属性后写入 DWD 明细表。DWS 层做轻度聚合。按业务主题把明细聚合成每分钟、每小时的汇总指标写入 Hologres 的聚合结果表。带宽敏感型指标比如 UV、GMV会在这里做精确去重。ADS 层面向具体应用比如实时大屏、实时风控规则引擎、运营后台看板。ADS 层数据基本都是从 DWS 聚合表里查询或者在 DWS 之上再套一层宽表视图。这套分层的核心价值和离线数仓一样隔离变更风险。上游表结构发生变化时DWD 层挡一道下游业务逻辑频繁变时DWS/ADS 层挡一道。但实时链路里每一层都会增加延迟所以不是所有表都需要四层齐全。我在实际项目里遵循一个原则核心交易链路三层保底ODS→DWD→DWS/ADS日志分析链路两层就够ODS→DWS别为了分层而分层。2.2 Hologres 表模型选型行存、列存还是行列共存Hologres 建表时第一个要决定的就是存储模式这个选错后面改起来非常痛苦。我整理了三种模式的适用场景。存储模式适合场景典型写入方式典型查询场景行存高并发点查、KV 查询Flink 主键 Upsert根据主键查单行或多行列存大范围聚合、OLAP 报表Flink 批量导入GROUP BY 的报表查询行列共存既要点查又要聚合分析Flink Upsert 写入点查 聚合混合我在实际建设中的经验是ODS/DWD 明细表基本用列存因为明细表的主要消费方式是按维度字段做聚合且数据量很大列存压缩比高、扫描效率好只有极少数需要按主键频繁点查的明细比如“查某个订单的最新状态”才考虑行存或行列共存。DWS 层的轻度聚合表反而适合行列共存——前端看板可能按精确 ID 查询某条指标同时也要按时间范围拉多天的趋势曲线。另外要注意 Hologres 的分区表和分布键设计。分布键决定了数据在哪个 Shard 上选择分布键时尽量和查询的等值过滤条件字段对齐这样能减少 Shuffle。Flink 写入时如果表的分布键和主键一致Upsert 的效率会更高因为同一主键的数据一定落在同一个 Shard 内更新时不需要跨节点协调。2.3 网络、账号与权限最容易被忽视的前置步骤很多人建实时链路时一上来就写 Flink SQL结果跑起来全是连接超时或者权限报错。这里有几个必须提前确认的事项。第一网络连通性。Flink 作业运行所在的 VPC 必须能访问到 Hologres 实例的 Endpoint。跨 VPC 要用对等连接或公网地址但公网访问会增加延迟且不稳定生产环境强烈建议走内网。Hologres 的 Endpoint 分 OLAP 地址和 FE 地址Flink 连接器用的是 FE 地址。第二账号权限。Hologres 的鉴权体系兼容 PostgreSQL建议为 Flink 任务单独创建一个专用账号只授予需要的库表权限。注意 Hologres 的权限模型里建表、写入、删除、DML、DDL 的权限是分开的Flink 写入时如果表不存在且开启了自动建表需要在账号上额外授予建表权限。第三关于开发调试我在测试环境喜欢用holo_admin这类超级账号快速验证但上了生产必须换成最小权限账号。否则一旦 Flink 任务有 Bug可能连带把整库的元数据搞乱。3. 核心玩法Flink 连接器接入 Hologres 的完整配置3.1 版本匹配与依赖包选择Flink 连接 Hologres 用的是官方提供的 holo-connector不同连接器版本对应不同 Flink 大版本比如 Flink 1.13、1.15、1.17 各有对应的 connector 包。这里最容易踩的坑是 Flink 小版本和连接器小版本不一致导致 NoSuchMethodError 或 ClassNotFoundException。我一般建议先用阿里云官方文档里的版本映射表确认然后直接用对应版本号的 connector。如果用的是 DataWorks 或 FDB 这类托管平台平台里通常已经内置了 holo-connector直接在 SQL 作业里声明connector hologres即可不需要自己上传 JAR如果是自建 Flink 集群需要把 connector JAR 放到 Flink 的lib目录或者在提交作业时通过-j参数显式指定。另一个容易忽略的是 JDBC Driver 版本。Hologres 的驱动分为 V1 和 V2 两代V2 驱动性能更好但有些老版本的 holo-connector 不支持。如果你的 Flink 任务写入吞吐上不去检查一下是不是连接器自带的驱动版本太旧。3.2 Flink SQL 创建 Hologres 结果表的标准 DDLFlink SQL 写 Hologres 的表分四种情况结果表Sink、源表Source、维表Lookup、数据变更订阅Binlog。最常用的是结果表和维表。先看一张典型的结果表 DDLCREATE TABLE holo_dwd_order_detail ( order_id BIGINT, user_id BIGINT, product_id BIGINT, amount DECIMAL(12, 2), order_status STRING, order_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector hologres, jdbcUrl jdbc:hologres://endpoint:port/dbname, dbname dbname, username username, password password, sink.batch.max-row-size 4096, sink.batch.write-time 10, sink.max-retry-times 3 );这里有几个关键点。PRIMARY KEY (order_id) NOT ENFORCED不是摆设Hologres 连接器会读这个字段来推断 Upsert 逻辑。如果源表没声明主键连接器会退化成 Insert 模式重复数据会累积声明了主键Flink 会把相同主键的多条数据合并后写入Hologres 内部执行存在则更新、不存在则插入这样 Flink Checkpoint 重放时不会产生重复记录。NOT ENFORCED的意思是 Flink 自己不校验主键唯一性主键约束由下游 Hologres 来保障。如果 Hologres 表本身也定义了主键两者必须完全一致否则写入时会报错。3.3 连接器参数逐项拆解与调优建议连接器参数是性能调优的主要抓手。我按写入链路的处理顺序整理了一张参数表方便对照检查。参数名作用建议值sink.batch.max-row-size单批次最大行数达到后触发写入4096~16384视单行数据大小调整sink.batch.max-size单批次最大字节数达到后触发写入20MB~50MBsink.batch.write-time最大攒批时间不管攒够多少都触发写入10~30 秒延迟敏感场景设小sink.max-retry-times单批次写入失败后的最大重试次数3~5 次sink.buffer-flush.pending-size缓冲区待写入数据大小与 batch.max-size 配合connectionSize连接池大小默认 3高并发写入可调到 5~8sink.ignore-write-exception是否忽略写入异常继续处理生产环境建议 false攒批参数是整个写入性能的命根子。Hologres 是分布式系统每次 RPC 有固定开销如果 Flink 一条一条写吞吐必然上不去反过来如果攒批时间太长数据延迟又会变大。我通常会先给一个保守配置max-row-size4096write-time10然后看线上监控的“单批次写入耗时”和“攒批耗时”占比来微调。如果攒批耗时占比高说明批次太小加大行数如果写入耗时占比高说明 Hologres 侧压力大优先看实例规格和连接池而不是继续加大批次。维表的 DDL 写法略有不同需要显式声明为维表CREATE TABLE holo_dim_user ( uid BIGINT PRIMARY KEY NOT ENFORCED, user_name STRING, user_level STRING, city STRING ) WITH ( connector hologres, jdbcUrl jdbc:hologres://endpoint:port/dbname, dbname dbname, username username, password password, lookup.cache.max-rows 10000, lookup.cache.ttl 300s );维表的缓存参数特别重要。lookup.cache.max-rows控制最多缓存多少行维度数据lookup.cache.ttl控制缓存时间。如果不加缓存每条主表数据都要触发一次维表查询大流量下 Hologres 会被点查打爆但缓存开太大又会让维度数据更新不及时我一般控制在 300 秒 TTL、1 万行以内具体看维表变更频率。4. 实战 SQL 拆解CDC 入仓、实时 ETL 与维表 JOIN4.1 用 Flink CDC 采集 MySQL 业务库实时数仓的上游通常不止 Kafka还有业务库的变更数据。Flink CDC 是采集 MySQL/PostgreSQL 增量数据的标准做法。我的习惯是直接读 Kafka 中的业务 Topic 并把消息体解析进 ODS 层但如果是自建 MySQL 没有成熟的消息管道用 Flink CDC 就能省掉 Kafka 这一跳。CREATE TABLE mysql_users_source ( id BIGINT, user_name STRING, email STRING, created_at TIMESTAMP(3), update_at TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 10.0.x.x, port 3306, username cdc_user, password ********, database-name app_db, table-name users, scan.startup.mode latest-offset ) ;scan.startup.mode有两个常用取值需要注意。initial会先做全量快照再增量同步适合第一次建链路时数据落底latest-offset只从当前位点开始读取增量适合已有全量数据做初始化、业务只关心新数据的场景。我建 ODS 层时通常先用initial跑一遍然后停掉作业改成latest-offset避免每天重复读全量。这里多说一句 CDC 的潜在风险。MySQL CDC 底层要解析 Binlog如果业务库开启了binlog_row_imageFULL连接器才能拿到完整的前后镜像否则更新事件里只有变更字段导致某些大宽表 JOIN 时字段丢失。建链路前最好找 DBA 确认一下 Binlog 配置否则等上线后再发现就非常被动。4.2 实时 ETL 写入 Hologres一条完整的 INSERT INTO 语句拿到源表数据后实时 ETL 的逻辑基本上就是“过滤 字段清洗 维表补全 分区裁剪”。下面是一条比较完整的 SQL演示了从 MySQL CDC 源表 JOIN Hologres 维表后写入 Hologres DWD 层的过程INSERT INTO holo_dwd_user_order ( order_id, user_id, user_name, city, amount, order_status, order_time, etl_time ) SELECT o.order_id, o.user_id, u.user_name, u.city, o.amount, o.order_status, o.order_time, NOW() FROM mysql_orders_source AS o LEFT JOIN holo_dim_user FOR SYSTEM_TIME AS OF o.proctime AS u ON o.user_id u.uid WHERE o.amount 0 AND o.order_status NOT IN (CANCELLED_INIT, INVALID);这段 SQL 对新手来说最有价值的是FOR SYSTEM_TIME AS OF o.proctime这个语法。它指定了维表 JOIN 的快照时间语义意思是每条源表数据在进入 Flink 时取那个时刻的维表数据进行关联这保证了流处理中维表数据不会因为更新而出现“先看到新值、后看到旧值”的混乱。需要注意的是源表 DDL 里如果要使用proctime必须显式声明处理时间字段CREATE TABLE mysql_orders_source ( ... order_time TIMESTAMP(3), proctime AS PROCTIME() ) WITH (...);ETL 时间我习惯用NOW()直接打点这样排查数据延迟时能明确知道这条记录是 Flink 什么时候处理完的。4.3 DWS 层实时聚合与 Hologres Binlog 回读DWS 层的典型指标是“最近 5 分钟 GMV”“今日 UV”“分城市订单量”。Flink SQL 里做这种聚合非常简单CREATE TABLE holo_dws_city_gmv ( stat_date STRING, city STRING, gmv_5min DECIMAL(12, 2), order_cnt BIGINT, window_time TIMESTAMP(3), PRIMARY KEY (stat_date, city, window_time) NOT ENFORCED ) WITH (...); INSERT INTO holo_dws_city_gmv SELECT DATE_FORMAT(window_start, yyyy-MM-dd), city, SUM(amount), COUNT(DISTINCT order_id), window_start FROM TABLE(TUMBLE(TABLE source_with_city, DESCRIPTOR(order_time), INTERVAL 5 MINUTE)) GROUP BY window_start, city;这里有一个真正的生产要点不要在 Flink 里用精确去重处理超大基数指标。COUNT(DISTINCT order_id)在数据量大的时候会吃掉大量内存Flink 的 State 会越来越大。我的替代方案是明细层把order_id落库DWS 层只做 VOID 方案写入明细表 用 Hologres 的近似计数或精确去重查询兜底。如果非要实时精确去重需要用 Bitmap 或 RoaringBitmap 做方案设计这部分复杂度较高建议先量化业务对“精确”的容忍度再决定。还有一个我自己很常用的技巧Hologres 表开启 Binlog 功能后Flink 可以用 holo-connector 把 Hologres 当作源表来读实现“实时数仓回读”。典型场景是DWS 表的数据由 Flink 写入另一些指标比如跨天累计量需要基于“当前累计 增量日志”计算那就直接用 Binlog 订阅 DWS 表的变更再做一层实时后处理。这个能力把 Hologres 从一个静态查询引擎变成了一个可以继续流转的中间态存储对复杂指标体系很有帮助。5. 性能与稳定性写入瓶颈、任务调参与一致性保障5.1 性能瓶颈的定位方法先看 Flink 还是先看 Hologres实时链路性能出问题第一反应不该是加并行度而是先判断瓶颈到底在 Flink 还是 Hologres。我给团队定了一个排查顺序。先看 Flink 界面上每个算子节点的 BackPressure 状态。如果 Source 到 Sink 之前的某个算子出现高背压说明计算逻辑或者 JOIN 逻辑卡住了跟 Hologres 关系不大如果背压集中在 Sink 算子说明写入端是瓶颈。此时看两个指标Sink 算子的currentLowWatermark是否持续不前进——如果 LowWatermark 不前进说明数据攒在 Flink 内部处理不了再对比 Hologres 侧监控里的 TPS 和响应时间如果 Hologres 的 QPS 不高但响应时间长大概率是连接池不够或者出现了锁等待。我之前遇到过一个问题Fl集群有 16 个并行度Hologres 实例规格也足够但写入速率就是上不去。后来排查发现 holo-connector 的默认连接池大小只有 316 个并行度抢 3 个连接大部分时间都耗在等待连接上。把connectionSize调到 8 之后吞吐直接翻了两倍。5.2 写入参数、并行度与分区键的协同调整调参不是单点调优而是整套参数的协同。我总结了一套组合拳按顺序执行效率最高。首先确认 Hologres 目标表的分区键和 Flink 写入并行度之间的关系。如果 Flink 的并行度远高于目标表的分区数写入时会产生大量跨 Shard 的分布式事务严重影响性能。一般建议 Flink Sink 并行度不要超过目标表分片数的 3~4 倍超过就先调 Sink 并行度而不是盲目堆资源。其次是攒批参数。数据源本身的流量模式也要考虑进去如果上游流量是突发型的比如大促峰值秒级暴涨固定大小的攒批策略不够灵活。我会把sink.batch.write-time设成一个中间值5~10 秒既保证低流量时的延迟又能在高流量时快速攒满批次。对于极低流量的表比如每天只有几百条变更加缓存反而会延迟这类表可以直接关闭攒批或者把write-time调小到 1 秒。最后是 Hologres 侧的写入引擎参数。如果实例规格是标准型写入吞吐不够时优先扩容 Shard而不是堆 CPU。Hologres 对写入吞吐的设计思路是“先扩容 Shard 再调并发”顺序反了的话加连接数也没用。5.3 一致性保障Checkpoint 语义、幂等写入与脏数据隔离实时链路的数据一致性是老大难我的方案是三层配合。Flink 侧开启 Checkpoint这是基础。exactly-once模式加上 Hologres 的连接器支持能在作业重启时保证端到端不重不丢。这里要注意Checkpoint 间隔不是越小越好间隔太短会导致频繁做快照占用大量带宽和磁盘 IO。我一般设置 60~120 秒在恢复时间和稳定性之间找平衡。Hologres 侧靠主键 Upsert 兜底。只要目标表有主键Flink 重启后即使重复输出了某些数据Hologres 也会用新数据覆盖旧数据。所以我在设计结果表时能定义主键的尽量定义主键主键字段要选业务上真正唯一的一组字段。如果表没有唯一业务键那就接受“至少一次”的语义在查询侧用时间去重。脏数据隔离是很多人忽略的环节。我在每张 DWD 结果表旁边都会建一张以_err结尾的异常表Flink 源表 DDL 里加上WITH (filter.invalid true)或者用WHERE条件把格式非法的数据单独 INSERT 到异常表里。这样上游数据格式变化时不会因为一条脏数据导致整个作业失败重启排查问题时异常表也直接给出了样例。6. 生产环境踩坑实录与完整排查链路6.1 坑一Flink JDBC 连接器异常作业频繁重启这个坑是在线任务上线第二周遇到的现象是作业运行几小时后突然报错Caused by: java.sql.SQLTransientConnectionException: Connection is not available, request timed out。排查链路如下。第一步看日志确认是哪个算子报错。发现集中在 Sink 节点排除源端和计算逻辑问题。第二步看 Flink 监控里的连接池指标和 Hologres 侧活跃连接数。发现 Hologres 的活跃连接数一直顶满配额而 Flink 侧并没有释放连接。第三步查连接泄漏点——最终定位到是某个并行度过高的作业同时打开了几百个连接加上 holo-connector 的连接回收机制对某些超时场景处理不及时导致连接被占满。处理方案是双管齐下调低该作业的 Sink 并行度同时把connectionSize从默认值改为 4又在 Hologres 侧调高了实例的最大连接数。之后这个坑没再出现过。这里有个经验连接池的大小不是越大越好连接数超过 Hologres 实例的规格上限后等待队列反而会放大延迟。6.2 坑二主键冲突与数据类型映射错误接手的第二个问题是数据写入报duplicate key value violates unique constraint。Hologres 主键表理论上支持 Upsert但报这个错说明连接器把数据当成普通 Insert 在提交。检查 Hologres 表的 DDL 发现目标表有主键但 Flink 端源表 DDL 里忘了写PRIMARY KEY ... NOT ENFORCED连接器认为下游表无主键于是走 Append 模式最终触发唯一约束冲突。解决方案是删掉 Flink 作业把结果表 DDL 补上主键声明重启。后来又遇到过数据类型映射问题Hologres 的DECIMAL(12, 2)在 Flink 里映射成了DECIMAL(20, 2)两边精度不一致导致写入报错。我的建议是建 Flink 结果表时统一用简化后的类型比如金额字段全链路都定义成DECIMAL(12, 2)不要在某个环节放大精度否则后面看数据时会发现一堆无意义的尾数。6.3 坑三CDC 表结构变更引发的全量重读这是最让我头痛的一个问题。MySQL 业务库有 DBA 在某个凌晨给users表加了个字段第二天早上发现 Flink CDC 作业直接失败日志里报 binlog 解析异常。根因是 CDC 连接器拿到的 Binlog 事件结构和它缓存的表结构不一致导致反序列化失败。解决方案分临时和长期两步。临时方案从最近一次 Checkpoint 恢复但因为是结构变更Checkpoint 里的状态已经无法匹配最终还是用initial模式重新跑了一轮全量 增量。长期方案和 DBA 约定任何上游表结构变更必须先通知数据团队在变更前停掉 CDC 作业变更后用latest-offset重启或者直接用支持 Schema Evolution 的 Canal DataHub 链路。这个坑让我养成了一个习惯所有实时链路的上游表只取业务上真正需要的字段不做SELECT *。这样即使上游加字段只要不加在已选字段上作业就不会挂。6.4 坑四Flink 火焰图排查 CPU 瓶颈有一次作业 CPU 使用率持续 90% 以上但背压并不严重数据吞吐也正常。这种反常现象用火焰图排查特别有效。Flink UI 的 JVM 指标里可以打开 CPU 采样火焰图能直观看到热点函数。那次火焰图显示耗时集中在 JSON 序列化上DataStream API 里每处理一条数据就调一次JSON.parseObject大促期间数据量放大后 CPU 直接被拖垮。解决方式是改用 Flink SQL 内置的JSON_VALUE和JSON_QUERY函数处理嵌套 JSON或者把解析逻辑放在源端提前完成避免在 Flink 内部重复解析。火焰图是定位运行时性能问题的最好工具之一但注意它只能看 CPU 热点内存问题、网络问题还是要配合 GC 日志和网络监控一起看。个人经验实时数仓最终还是工程问题把 Flink 和 Hologres 这条路跑通之后我的整体感受是技术选型只是开始真正的难点在工程化。连接器参数要迭代调优表结构设计要反复权衡上下游的变更管理要有规范的流程数据一致性和延迟之间的取舍要不断和业务对齐。实时数仓听起来很高大上落到地上就是一张表、一条 SQL、一次重启、一行日志积累出来的。最后分享一个我自己很受用的习惯每张 Hologres 结果表建表时把主键、模式、近一周的写入峰值、查询方式这些信息都维护起来。实时链路的问题基本都不在“跑不通”而在“跑着跑着就出事”做好元数据管理能让你在出事时快速定位而不是从头开始翻代码。希望这篇实践总结能帮你少走一些弯路。