简介:这份PDF文档面向数据工程师、架构师及实时数仓方向的技术人员,围绕Flink与Hologres的组合,讲解云原生环境下实时数仓的构建思路与落地实践,帮助读者理解如何解决传统数仓延迟高、架构复杂、资源消耗大等问题。内容涵盖Lambda架构的局限、HTAP与HSAP理念的演进、实时导入与批量归档、维表关联与离线加速、联邦计算、结果缓存、计算存储分离及云原生统一存储等关键主题,并结合Hologres兼容PG语法、BI对接、向量检索等能力展开分析。资源包共1个PDF文件,大小约1.23MB,便于下载后离线阅读与查阅。目前已有593人学习浏览,适合希望系统了解实时数仓选型、架构简化与流批一体方案的技术人员参考,可从中获取架构设计思路、技术选型依据与最佳实践要点。
1. 从 Lambda 到 HSAP:这份 Flink + Hologres 实践文档到底解决了什么问题
如果你正在维护一套 Hive + Flink + Impala + Kudu 的实时数仓,大概率遇到过这种场景:业务方要一个多维分析报表,你得同时维护离线链路和实时链路,两边数据对不上还得排查半天;大促期间写入吞吐一上来,Kudu 的写入延迟就开始抖,查询更是没法保证秒级响应。这份《Flink Hologres 云原生实时数仓最佳实践》文档,讲的就是怎么把上面这套架构收敛成一条链路——用 Flink 做流批统一的加工,用 Hologres 做统一存储和服务出口,把实时数仓和在线数据服务融合到一个引擎里。
文档的核心主张是 HSAP(Hybrid Serving & Analytical Processing),也就是分析和服务一体化。它跟 HTAP 的区别在于:HTAP 面向的是有事务需求的业务系统,需要保证 ACID 和 TP/AP 一致性;而 HSAP 面向的是埋点、机器数据这类高吞吐写入场景,不需要事务开销,重点在于统一实时和离线的存储引擎,同时支撑高 QPS 点查和 PB 级 OLAP 查询。适合谁看?正在做实时数仓选型、被 Lambda 架构的运维成本折磨、或者想把在线服务和离线分析统一到一个存储出口的团队。文档里给了两个真实客户案例——考拉和阿里 CCO 体验系统,有具体的性能数字和架构演进路径,不是纯理论。
2. HSAP 架构选型:为什么是 Hologres 而不是 ClickHouse 或 Druid
2.1 Lambda 架构的账算不过来在哪
Lambda 架构的经典问题是三条链路并行:离线数仓跑 T+1 的批处理,实时数仓跑 T+0 的流处理,中间还要一个融合层来对齐两边数据。文档里列了几个硬伤——架构复杂、资源消耗大、数据孤岛、人才培养难、开发成本高、不敏捷。这些不是空话,落到日常运维就是:同一份业务逻辑要在 Hive SQL 和 Flink SQL 里各写一遍,口径对不上就是血泪排查;实时链路和离线链路各占一套资源,大促时两边都要扩容,成本翻倍;新人进来得同时学 Hive、Flink、Impala、Kudu 四套东西,培养周期长。
文档提出的思路是实时离线一体化加分析服务一体化。实时离线一体化指的是统一存储引擎,实时写入的数据和批量导入的数据落在同一个存储里,不需要在融合层做对齐。分析服务一体化指的是同一个引擎既能跑 OLAP 复杂查询,又能扛高 QPS 点查,不需要在 ClickHouse 和 HBase 之间来回倒数据。
2.2 Hologres 的 HSAP 定位与关键能力
Hologres 是阿里自研的大数据体系组件,定位是 HSAP 引擎。文档里列了几个关键能力,我挑对选型最有影响的几个说:
兼容 PostgreSQL 语法。这意味着 PG 生态的开发运维工具可以直接用,BI 工具对接不需要额外适配层。对于已经熟悉 PG 的团队,上手成本几乎为零。
行列共存存储。列存对分析友好,行存对点查快速。Hologres 支持在同一张表里同时建行存和列存,或者根据查询模式自动选择。这个能力直接决定了它能不能同时扛 OLAP 和 Serving 两类负载。
计算存储分离。计算资源和存储资源独立扩缩容,按需使用。文档里提到与 MaxCompute 底层打通,可以透明加速,减少数据搬迁。这个对成本控制很关键——存储便宜、计算贵,分离之后可以只扩计算不扩存储。
C++ Native 执行引擎 + 优化器。向量化、全异步执行,轻量级用户态线程调度,同时支持高并发和复杂统计两类负载。公平调度算法(CFS)保证高并发场景下计算资源充分利用。
文档里给了一个架构对比表,我整理成更直观的形式:
| 维度 | Lambda 架构 | HSAP 架构 |
|---|---|---|
| 存储引擎 | 离线 Hive + 实时 Kudu/HBase | 统一 Hologres |
| 数据一致性 | 融合层对齐,T+1 和 T+0 可能不一致 | 写入即可见,单一数据源 |
| 查询类型 | OLAP 走 Impala/ClickHouse,点查走 HBase/Redis | OLAP + 点查统一引擎 |
| 扩缩容 | 离线/实时各自扩,成本高 | 计算存储分离,按需弹性 |
| 开发成本 | 同一逻辑写两遍,口径对齐耗时 | 一套 SQL,流批统一 |
2.3 从开源 Hadoop 迁移到云原生的实际路径
文档里考拉的案例给了具体迁移路径。原来的技术栈是 Hive + Flink + Impala + Kudu,业务诉求是:运维成本高、数据写入慢(不支持 100w+/s)、查询可见有延迟、数据量暴增下查询性能无法保证、离线近实时准实时多条链路共存架构冗余。
迁移后的架构是 MaxCompute + Hologres + Flink。几个关键设计决策:
数据服务单一出口在 Hologres。所有查询——多维分析、点查、报表——都走 Hologres,不再维护多个查询引擎。
CDM 层数据持久化。Common Data Model 层的数据落在 Hologres 里持久化,方便数据回刷和实时/离线差异排查。这个设计很实用——出问题时可以直接对比 CDM 层的数据,不用去翻 Kafka 消息。
聚合操作在 Hologres 进行。Flink 层只负责清洗和强指标计算,聚合操作下沉到 Hologres。这样降低了流处理任务的压力,也利用了 Hologres 的分层数据模型优势。
多维查询结果存储到 Hologres。Blink 层负责清洗和强指标计算,结果写入 Hologres 供查询。
迁移效果:几十亿商品的特征信息仅耗时 5 分钟完成数据切换;维度变更链路 1 小时内完成维表数据切换,无需更改 Flink 作业;支持自助即席多维分析,涵盖 1000+ 自定义维度信息。
3. Flink + Hologres 实时数仓搭建:从数据接入到分层建模
3.1 整体架构与数据流向
文档里的架构图拆解下来是这样的数据流:
数据源(RDS/日志/埋点) ↓ Flink 实时接入(Kafka/DataHub) ↓ Flink 清洗、关联、转换 ↓ Hologres DWD 明细层(列存) ↓ Hologres DWS 汇聚层(列存 + 行存) ↓ Hologres ADS 应用层(行存) ↓ OLAP 报表 / 点查服务 / 在线应用几个关键设计点:
DWD 层用列存。明细数据量大,列存压缩比高,扫描效率好。Flink 写入时直接写列存表。
DWS 层列存 + 行存混合。轻度汇总数据用列存供分析查询,高度汇总数据用行存供点查。文档里提到“明细数据列存、轻度汇总数据列存、高度汇总数据行存”的分层策略。
ADS 层用行存。面向点查、监控、在线类服务,行存对单行读取更快。
维表数据用行存。维度数据需要频繁关联,行存点查性能好。
3.2 Flink 实时写入 Hologres 的配置要点
Flink 写入 Hologres 常见做法是用 Hologres 提供的 Flink connector。下面是一个典型的 Flink SQL 建表语句,我按文档里的架构补全了参数:
-- Flink SQL 建表:写入 Hologres DWD 层 CREATE TABLE dwd_user_behavior ( user_id BIGINT, item_id BIGINT, behavior_type STRING, event_time TIMESTAMP(3), proc_time AS PROCTIME() ) WITH ( 'connector' = 'hologres', 'dbname' = 'realtime_dw', 'tablename' = 'dwd_user_behavior', 'username' = '${access_id}', 'password' = '${access_key}', 'endpoint' = '${hologres_endpoint}', 'connectionSize' = '10', -- 连接池大小,高吞吐场景适当调大 'jdbcWriteBatchSize' = '1024', -- 批量写入条数,默认 256,大促可调到 2048 'jdbcWriteFlushInterval' = '10000',-- 刷写间隔(ms),默认 10s 'mutateType' = 'insertorupdate', -- 支持更新,写入即可见 'partition' = 'ds' -- 分区键,按天分区 );参数说明:
connectionSize:连接池大小。写入吞吐上不去时优先调这个,但不要超过 Hologres 实例的并发上限。jdbcWriteBatchSize:批量写入条数。默认 256 偏保守,大促场景可以调到 1024 或 2048,但要注意单批次数据量不要超过 Hologres 的写入限制。jdbcWriteFlushInterval:刷写间隔。设太短会导致小文件多,设太长会导致数据可见延迟高。文档里 CCO 案例的写入延迟稳定在 500us 内,这个参数需要配合 batchSize 一起调。mutateType:insertorupdate支持主键更新,适合维表关联后的宽表写入;纯追加场景用insert性能更好。
3.3 维表关联与离线加速的实现方式
维表关联是实时数仓里最容易翻车的环节。文档里提到的做法是:维度数据存在 Hologres 行存表里,Flink 作业通过 JDBC 或 Hologres connector 做 lookup join。
-- Flink SQL:维表关联 CREATE TABLE dim_item ( item_id BIGINT, item_name STRING, category_id BIGINT, update_time TIMESTAMP(3), PRIMARY KEY (item_id) NOT ENFORCED ) WITH ( 'connector' = 'hologres', 'dbname' = 'dim_db', 'tablename' = 'dim_item', 'username' = '${access_id}', 'password' = '${access_key}', 'endpoint' = '${hologres_endpoint}', 'lookupCacheMaxRows' = '100000', -- 缓存行数,减少对 Hologres 的查询压力 'lookupCacheExpireTime' = '60000' -- 缓存过期时间(ms),维度变更后 1 分钟内生效 ); -- 关联查询 INSERT INTO dwd_user_behavior_wide SELECT b.user_id, b.item_id, d.item_name, d.category_id, b.behavior_type, b.event_time FROM dwd_user_behavior AS b LEFT JOIN dim_item FOR SYSTEM_TIME AS OF b.proc_time AS d ON b.item_id = d.item_id;这里的关键参数是lookupCacheMaxRows和lookupCacheExpireTime。文档里考拉案例提到“维度变更链路 1 小时内完成维表数据切换”,实际配置时缓存过期时间决定了维度变更的生效延迟。设太短会导致频繁查 Hologres,设太长会导致维度变更不及时。常见做法是设 60s 到 5 分钟之间,根据维度变更频率调整。
离线加速的思路是:对于历史数据的复杂查询,利用 Hologres 与 MaxCompute 的底层打通能力,直接查询 MaxCompute 外表,减少数据搬迁。文档里提到“MaxCompute 无缝打通,减少数据搬迁,透明加速”。
3.4 数仓分层建模的敏捷化实践
文档里给了一个分层建模的实践框架:
ODS 层:数据归集。原始数据从 DataHub 或 Kafka 接入,不做太多加工,保留原始字段。
DWD 层:数据加工。Flink 做清洗、关联、转换,输出明细宽表。这一层是实时数仓的核心,数据质量直接决定上层可用性。
DWS 层:多维分析、数据集市。轻度汇总,面向主题的、可共享的数仓分层建设。文档里强调“加工服务一体化:在 Flink 中加工,在 Hologres 中服务,减少数据移动,减少数据孤岛”。
ADS 层:报表类、服务化。高度汇总,面向具体应用场景。文档里提到“弱化 ADS、面向 DWS、DWD 的应用开发服务”,意思是尽量让应用直接查 DWS 层,减少 ADS 层的维护成本。
这个分层策略的核心思想是减少数据层次,敏捷适应需求变化。传统数仓 ODS→DWD→DWS→ADS 四层,每层都要单独维护 ETL 逻辑。Flink + Hologres 的方案里,Flink 负责 DWD 层的加工,DWS 和 ADS 层通过 Hologres 的视图或物化视图实现,减少数据搬迁。
4. 避坑与排查:Flink 写 Hologres 常见的五个翻车场景
4.1 写入吞吐上不去,Flink 作业反压
现象:Flink 作业的 Sink 算子反压指标持续飙红,Hologres 侧写入延迟升高,Kafka 消费 lag 增长。
原因:常见的有三种——连接池太小导致写入并发不足;批量写入条数设置过小导致频繁网络往返;Hologres 实例的 Shard 数不够,写入热点集中在少数 Shard 上。
解决:先看 Flink 的反压指标定位是 Sink 侧还是上游。如果是 Sink 侧,调大connectionSize和jdbcWriteBatchSize。如果 Hologres 侧 CPU 不高但写入延迟高,检查表的 Shard 数——Hologres 建表时可以指定shard_count,默认值可能不适合高吞吐场景。文档里 CCO 案例的双 11 峰值 TPS 输入 100w+/s,这种量级需要提前做好 Shard 规划。
4.2 维表关联数据不一致,维度变更后实时数据没更新
现象:维度表更新后,Flink 作业关联出来的宽表数据还是旧值,延迟很久才生效。
原因:lookupCacheExpireTime设置过长,或者维表关联用了FOR SYSTEM_TIME AS OF但缓存没有正确失效。
解决:检查lookupCacheExpireTime配置,根据维度变更频率调整。如果维度变更需要秒级生效,可以关闭缓存或者设一个很短的过期时间,但这样会增加 Hologres 的查询压力。折中方案是用 Hologres 的 Binlog 能力,维度变更时主动通知 Flink 刷新缓存。文档里考拉案例做到“维度变更链路 1 小时内完成维表数据切换”,这个延迟对于大多数场景够用。
4.3 Hologres 查询延迟高,OLAP 和点查互相影响
现象:跑一个复杂的 OLAP 查询时,在线点查的延迟从毫秒级飙升到秒级。
原因:OLAP 查询占用了大量计算资源,点查请求排队。Hologres 虽然有公平调度算法,但如果资源本身不够,调度也救不了。
解决:利用 Hologres 的计算组(Compute Group)能力,把 OLAP 查询和点查分配到不同的计算组,物理隔离资源。文档里提到“轻量级用户态线程调度,同时支持多种查询负载”,但实际生产环境还是建议做资源隔离。另外,点查场景尽量走行存表,OLAP 走列存表,避免互相争抢。
4.4 Flink 作业重启后数据重复写入
现象:Flink 作业从 Checkpoint 恢复后,Hologres 表里出现重复数据。
原因:Hologres Sink 的mutateType设成了insert,没有主键去重。或者 Checkpoint 间隔太长,恢复时重放了大量数据。
解决:如果业务允许更新,把mutateType改成insertorupdate,利用 Hologres 的主键做去重。如果必须是追加模式,需要在 Flink 侧做去重,比如用ROW_NUMBER()开窗去重,或者启用 Flink 的 Exactly-Once 语义(需要 Hologres Sink 支持两阶段提交)。文档里没有展开讲 Exactly-Once 的配置,但这是生产环境必须考虑的问题。
4.5 离线数据和实时数据对不上
现象:同一份数据,Hologres 实时查询的结果和 MaxCompute 离线查询的结果不一致。
原因:实时链路和离线链路的加工逻辑不一致,或者数据写入 Hologres 时有丢失。
解决:文档里考拉案例的做法是“CDM 层数据持久化,方便数据回刷以及实时/离线差异问题排查”。具体操作是:在 Hologres 里保留 CDM 层的明细数据,离线链路和实时链路都从 CDM 层开始加工,这样出问题时可以直接对比 CDM 层的数据,定位是接入问题还是加工问题。另外,Flink 作业的 Checkpoint 和 Hologres 的写入事务要配合好,确保数据不丢不重。
5. 从 CCO 双 11 案例看性能调优的边界与验证方法
文档里阿里 CCO 体验系统的案例给了很具体的性能数字,我拿它当基准来聊调优的边界。2019 年双 11 当天,Flink 实时数据加工峰值 TPS 输入 100w+/s,写入延迟稳定在 500us 内;MC-Hologres 查询服务当天查询 latency 平均 142ms,99.99% 的查询在 200ms 以内;支撑 200+ 实时数据大屏,为近 300+ 小二提供数据查询服务;同时支撑多维分析和高 QPS 服务化查询场景。整体硬件资源成本下降 60+%。
这些数字背后有几个调优动作值得拆开看。
写入延迟 500us 内怎么做到的。Flink 侧用了批量写入加异步刷写,Hologres 侧用了写友好的数据结构,支持高吞吐写入。文档里提到“写友好数据结构,高吞吐数据写入,支持更新,写入即可见”。实际配置时,jdbcWriteBatchSize和jdbcWriteFlushInterval的配合很关键——批量条数要足够大以减少网络往返,刷写间隔要足够短以保证可见性。500us 的延迟意味着刷写间隔可能在毫秒级,这对 Hologres 的写入能力要求很高。
99.99% 查询在 200ms 以内怎么验证。Hologres 提供了查询日志和慢查询分析功能。我一般会做三件事:第一,在 Hologres 侧开启慢查询日志,设置阈值比如 100ms,定期分析慢查询模式;第二,在 Flink 侧监控 Sink 的写入延迟和反压指标,确保写入不会成为瓶颈;第三,用压测工具模拟高并发点查和 OLAP 混合负载,观察 P99 和 P999 延迟。文档里没有给具体的压测方法,但这是上线前必须做的验证。
成本下降 60+% 的来源。计算存储分离是主要贡献——存储用便宜的 OSS 或 Pangu,计算按需扩缩容。另外,统一存储减少了数据搬迁和冗余存储。文档里提到“与 MaxCompute 底层打通,透明加速”,这意味着历史数据可以留在 MaxCompute 里,Hologres 只存热数据,进一步降低成本。
一个具体的验证技巧:在 Hologres 里建两张表,一张行存一张列存,写入相同的数据,然后分别跑点查和 OLAP 查询,对比延迟和资源消耗。这个测试能帮你确定业务场景下行存和列存的比例。我一般会建议点查为主的表用行存,分析为主的表用列存,混合场景用行列共存但要注意存储成本。
从那以后我每次做实时数仓选型,都会先跑一遍这个行存/列存对比测试,再根据结果决定分层策略。希望帮到你。
本文还有配套的精品资源,点击获取