搞实时数仓和OLAP分析平台这几年,我试过的方案不在少数,从最早的离线T+1跑批,到中间用Kafka+Flink自己拼装一套实时链路,再到后来在云上把DataWorks和Hologres组合起来做企业级实时数仓,这一路踩过的坑、填过的洞确实不少。今天这篇文章不聊概念,就把这套组合真正落地时的架构设计、表模型选择、查询优化、写入调优、运维成本这些实操层面的东西完整拆一遍,希望能给正在从离线数仓往实时化演进、或者正准备上OLAP分析平台的团队一些可以直接参考的经验。
先交代一下场景背景。假设你所在的公司已经有一定数据体量:几十个业务库、核心业务表每天新增几千万行、报表要求分钟级刷新、运营要看实时的漏斗和转化、管理层要看经营大盘。这个时候离线数仓的T+1延迟已经明显撑不住业务诉求,你需要一套能将实时链路和离线链路统一管理的平台,还要能承担高并发、低延迟的OLAP查询。DataWorks+Hologres的组合正好覆盖了这三个核心问题:数据开发与调度、数据治理、实时交互式分析。
1. 实时数仓为什么不能只靠“Flink单兵作战”
很多团队的第一反应是,实时数仓嘛,上Flink就完事儿了。Kafka里接数据,Flink做ETL,计算结果写回消息队列或者MySQL、Redis,业务方自己去消费。说实话,这种方案在小规模、单业务场景下是能跑的,但一旦要上升到“企业级”,你会发现纯Flink方案有几个绕不过去的短板。
第一个短板是OLAP查询能力基本靠外部系统补。Flink本质是流式计算引擎,它擅长的是无界数据的有状态计算,你想让它直接支撑多维分析、即席查询、BI报表,体验会很差。最终数据还是得落到某个存储引擎里,那这个引擎的查询能力、并发能力、SQL兼容性就变成了整个链路的瓶颈。
第二个短板是开发链路过长、成本高。实时任务要自己维护状态、管理Checkpoint、处理数据回溯;实时结果表要自己定义存储结构;指标口径要自己在多处代码里维护。业务方今天加一个维度、明天改一个指标,Flink任务每次都要重新上线,长此以往,实时数仓就变成了一堆谁也改不动、也不敢动的“黑盒任务”。
第三个短板是缺乏全链路的数据治理能力。离线数仓好歹有血缘、有数据地图、有质量监控,到了实时链路,很多团队就直接裸奔了。数据从哪来、经过哪些加工、产出哪些指标、谁在用这张表,完全讲不清楚,出了问题只能看日志硬排。
DataWorks+Hologres这套组合的定位,其实就是把“实时链路”拉回到“工程化治理”的轨道上。DataWorks负责统一的数据集成、任务开发、调度运维、数据质量、血缘管理;Hologres负责实时的数据存储和高性能OLAP查询。流计算的任务还是Flink在跑,但Flink只作为计算引擎,任务由DataWorks统一编排和发布,结果落到Hologres,由Hologres对外提供查询服务。这样既保留了Flink的实时计算能力,又把实时数仓纳入了企业级数据治理的范畴。
我在实际项目中见过太多“Flink一把梭”后留下的烂摊子:状态后端频繁膨胀、数据重复、指标口径失控、任务依赖关系完全不可见。所以我的建议很直接:实时数仓这件事,计算引擎可以选型,但平台和存储一定要选那种“能让你看清楚全局”的组合。DataWorks+Hologres最大的价值,不是单点性能跑得多快,而是把开发和治理闭环打通了。
2. Hologres的存储与查询引擎:列存、索引与分布式调度
既然Hologres在这套方案里扛的是“存储+OLAP”这个核心角色,那它的底层逻辑就值得先讲透。Hologres是阿里云自研的实时交互式分析引擎,它本质上是分布式Share Nothing架构,底层存储使用列存格式,支持行列共存,查询引擎采用向量化执行。
为什么“列存”对OLAP这么重要?因为OLAP查询往往是“少数列、大量行”的扫描模式。比如一张订单表有80个字段,业务方统计“昨天的销售额按城市分组”,真正参与计算的可能只有城市、金额、时间三个字段。列存可以把这三列连续读出来,IO量大大减少,再配合压缩,整列扫描的性价比远高于行存。
但Hologres没有走纯列存的极端路线,而是支持行存、列存、行列共存三种表存储模式。这一点我认为非常务实,因为真实场景里,一张表不只是被分析引擎读,还可能被业务系统或者Flink频繁做点查。纯列存虽然分析快,但单行更新、按主键查询的性能通常不如行存。行列共存的意思是一份数据同时维护行存和列存两套物理形态,读取时按查询类型自动选择,写入时同步更新。代价是存储成本和写入开销更高,所以我的经验是:只有那种“既要被点查、又要被分析”的核心维度表才值得用行列共存,普通明细表用列存就够了。
为了理解Hologres的数据分布方式,必须先搞清楚三个关键属性:分布键(distribution_key)、分区键(partition_key)、分段键(segment_key)。
分布键决定数据在各个Shard上的分布策略。Hologres建表时如果不指定分布键,默认会把所有列作为分布键做Hash分布;更常规的做法是明确指定一个业务上天然分散的字段,比如用户ID、设备ID。分布键选得好不好,直接影响Join和Group By是否有数据重分布开销。
分区键的作用是物理切分数据,最典型的用法是按时间分区,比如ds TEXT作为分区列。按天分区后,查询可以走分区裁剪,只扫描命中的分区,数据管理(删除过期数据)也方便。
分段键则是Hologres特有的设计,它控制数据在Shard内部的排列方式。比如订单流水表可以按成交时间作为分段键,这样同一时间段的数据在物理上连续存放,时间范围查询的扫描量会大幅减少。
我记得有一次优化一个订单明细查询,查询条件是“某用户在最近30天的订单”,数据量大概有10亿行。最初这张表的分布键没指定,默认全列Hash,查询时Shard内部也没有按时间聚簇,每次查询都要扫大量无关数据,响应时间一直在秒级徘徊。后来重建表,分布键指定为uid,分段键指定为order_time,同样的查询直接降到几十毫秒。这个案例足以说明,Hologres的表属性设计不是随便填填就行,它直接决定查询能不能走索引、能不能做裁剪、能不能减少网络Shuffle。
Hologres的查询引擎走的是全向量化执行,SQL兼容PostgreSQL生态,支持标准SQL和大部分PostgreSQL函数,这使得对接BI工具的门槛非常低。和ClickHouse相比,Hologres在SQL兼容性、数据更新能力、高并发点查上更均衡;和Doris相比,Hologres和Flink的深度集成做得更顺滑,实时写入链路的代码量更少。这也是我为什么在实际项目中更倾向于用Hologres做“实时数仓的主体存储”,而不是单纯追求极限扫描性能的ClickHouse。
-- Hologres建表示例(订单明细表) CREATE TABLE dwd_order_detail ( uid TEXT NOT NULL, order_id TEXT NOT NULL, city TEXT, amount DOUBLE PRECISION, order_time TIMESTAMPTZ, ds TEXT NOT NULL ); CALL SET_TABLE_PROPERTY('dwd_order_detail', 'distribution_key', 'uid'); CALL SET_TABLE_PROPERTY('dwd_order_detail', 'partition_key', 'ds'); CALL SET_TABLE_PROPERTY('dwd_order_detail', 'segment_key', 'order_time'); CALL SET_TABLE_PROPERTY('dwd_order_detail', 'clustering_key', 'order_time');上面clustering_key是聚簇索引,配合分段键,可以让相同业务时间的记录物理相邻。实际效果就是“点查+范围查+聚合查”的混合负载能同时跑得不错,这在实时数仓的明细层(DWD)非常实用。
3. DataWorks的职责边界:数据集成、任务编排与治理闭环
很多人会把DataWorks理解成“一个写SQL的网页工具”,这个认知太窄了。DataWorks在实时数仓体系里的核心价值体现在四个层面:数据集成、任务开发与调度、数据质量、数据治理。
先看数据集成。DataWorks的数据集成模块支持离线同步和实时同步两大类。实时同步可以直连MySQL等业务库的Binlog,也可以从Kafka订阅消息,实时写入Hologres。这一点非常关键,因为很多公司的实时链路是“全手工拼装”的,Binlog采集写一个程序、Kafka到Flink再写一个程序、Flink到Hologres再写一堆Connector配置,全链路的状态都停留在个人手里。DataWorks把实时同步做成可视化配置,谁建的、同步位点到哪了、延迟多少毫秒,一目了然。
在数据开发层面,DataWorks支持SQL任务、Shell任务、Flink任务等多种类型。实时数仓场景下,通常的模式是:用SQL任务写离线加工逻辑,用Flink任务写实时加工逻辑,所有任务在DataWorks上统一编排,配置依赖关系和调度周期。这套体系最大的优势是“实时任务和离线任务在一个平台维护”,统一的代码版本、统一的发布流程、统一的运行日志,而不是实时任务散落在个人电脑上跑。
调度运维这块,DataWorks的周期调度支持分钟级、小时级、天级任务,并支持实例维度的重跑、补数据、置成功、冻结等操作。实时数仓虽然强调“实时”,但你依然需要离线任务来对账、补数、修正历史数据,所以“调度平台同时管理流和批”这件事,在实战中远比想象中更重要。
数据治理闭环是我最看重的部分。DataWorks的数据地图可以自动采集Hologres表的元数据,生成字段级血缘关系。比如一个实时任务从Kafka读取、经过Flink加工、写入Hologres的DWS表,再被Quick BI报表读取,整条链路在数据地图上都能看到。另一个实用能力是数据质量监控,可以对Hologres表设置主键监控、表行数波动监控、字段空值率监控,一旦实时写入出现异常,立刻报警并触发关联任务通知。这个能力在实时链路里尤其重要——离线的天级任务有问题第二天能发现,实时链路数据错了一小时,业务方可能已经基于错误数据做了决策。
DataWorks和Hologres之间的联动还有几个顺手的功能值得提:DataWorks可以直接把Hologres表元数据批量导入到数据地图,也可以一键生成Hologres表的建表语句;DataWorks的数据服务可以把Hologres的查询能力封装成API,供业务系统调用。也就是说,从数据接入、数据开发、数据调度、数据质量到数据服务,这一整套东西都收敛在DataWorks这个平台上,而不是靠多个开源组件拼拼凑凑。对于需要“企业级”三个字的团队来说,这种收敛本身就是巨大的运维红利。
4. 分层数仓怎么搭:ODS到ADS的实时化改造实践
接下来是数仓架构设计的核心问题。很多人以为实时数仓就是“一张大宽表,Flink算完直接写进去,BI一查完事”,但稍微复杂一点的业务,这种粗暴做法很快就会出问题:指标口径混乱、存储冗余、需求变更成本高。真正的企业级实时数仓,仍然需要分层设计,只是每一层的实现方式和离线数仓有所区别。
我把分层架构归纳成下表,这是我在实际项目中反复调整后沉淀下来的做法:
| 分层 | 主要职责 | 数据来源 | 写入方式 | Hologres存储模式 |
|---|---|---|---|---|
| ODS | 原始明细、存原始日志/业务数据 | Kafka、Binlog、离线同步 | Flink实时写入 / DataWorks实时同步 | 列存(可按天分区) |
| DWD | 清洗、标准化、维度补充后的明细事实 | ODS层 | Flink SQL / DataWorks SQL任务 | 列存,设置分布键和分段键 |
| DWS | 按主题域汇总的轻度聚合结果 | DWD层 | Flink SQL 或周期调度SQL | 列存,按维度分桶/分区 |
| ADS | 面向具体业务场景的应用层数据 | DWS / DWD层 | SQL任务 / 数据服务API | 行存或行列共存,高并发查询 |
| DIM | 维度数据,用户、商品、城市等 | 业务库离线同步/实时同步 | DataWorks实时同步 | 行列共存,支持点查与关联查询 |
ODS层最核心的原则是“不改数据、留存原始”。实时链路中,ODS层一般直接从Kafka消费Binlog或日志数据,Flink只做格式转换,不做过多的业务清洗,然后写入Hologres的ODS表。Hologres本身支持主键更新,所以Binlog流里的Update和Delete操作也能正确处理,这一点比把数据单纯堆进Kafka再另外建索引要省事得多。ODS表的数据量通常最大,按天分区,配合Hologres的生命周期管理定期清理过期分区。
DWD层的建设是实时数仓里工程量最大的部分。以交易场景为例,DWD层要把订单流、支付流、退款流、物流流等事实数据做关联和标准化,同时补齐维度字段,比如把商品ID关联出类目、把店铺ID关联出商家名称。离线数仓的关联可以用天级任务慢慢跑,实时数仓则要依赖Flink的Join能力——比如用Flink的流表Join维表,从Hologres的DIM表里读取维度数据实时补全字段,然后写出到DWD表。这里有一个关键建议:凡是需要频繁关联的维度数据,一定要提前放到Hologres的行列共存表里,让Flink的维表Join走Hologres的异步Lookup,避免每条流记录都打一次同步查询。
DWS层适合放“轻度汇总”数据,比如“用户当日订单数”“商品当日销售额”“店铺30日累计GMV”。这层的数据用Flink SQL做滚动窗口或滑动窗口聚合,直接写入Hologres的DWS表。分布键一般为维度ID(用户ID、商品ID、店铺ID),查询时命中的Shard非常收敛,聚合性能极高。DWS表也是BI报表主要查询的表,因为数据量比明细小一个数量级以上,查询响应速度更容易做到秒级甚至毫秒级。
ADS层则完全面向场景,比如“大促实时战报”“运营实时漏斗”“流量实时看板”。这层表通常由SQL任务或服务API加工,查询并发高、Query模式固定。ADS表建表时建议仔细评估点查和分析的比例,点查多就选择行存或行列共存,纯分析型报表用列存即可。
这套分层架构里有一个容易被忽略的设计点:每一层都要考虑“实时数据和离线数据的双写兼容”。因为实时链路偶尔要出问题,你总需要一个离线批任务回补历史数据或者重算当天数据。我的做法是,DWD和DWS层的Hologres表在设计上同时兼容Flink实时写入和DataWorks离线SQL写入,两边的表结构完全一致,写入方式不同而已。出故障时,用离线重算覆盖实时数据,表结构不用动,业务查询层感知不到。这就是流批一体在Hologres上的典型落地方式。
5. OLAP查询性能实战:SQL优化、写入调优与BI接入
架构搭好了,接下来就是考验真实性能的地方。Hologres跑OLAP查询的性能,一半取决于表设计,另一半取决于SQL写法。把“查询慢”甩锅给引擎,在Hologres这里大概率是冤枉了它。
先说查询侧的几个高频优化手段。
分布键过滤要写进SQL。Hologres的查询如果能在SQL里带上分布键等值条件,引擎可以直接定位到对应的单个Shard,跳过大量数据的重分布开销。比如DWS表分布键是uid,查询SELECT ... WHERE uid = 'user123',这条查询只会落在对应Shard上执行,性能直接起飞。怕就怕业务方的报表工具生成的SQL不带分布键条件,那就只能在SQL里做子查询或者加一层视图来约束过滤条件。
分区裁剪必须生效。查询里尽量带上分区列的范围过滤,比如WHERE ds BETWEEN '20240101' AND '20240107'。Hologres的分区裁剪是优化器自动完成的,但前提是你得让优化器看到分区条件。如果SQL写成函数套分区列,比如WHERE date(ds) = '2024-01-01',裁剪往往会失效。这是我见过的高频踩坑点,在写SQL模板给BI团队复用的时候,务必把分区条件的写法定死。
聚合类SQL别整列扫描。统计类查询尽量只SELECT需要的列,别写SELECT COUNT(*) FROM 大表这种全表扫描的语句。建表时把常用过滤条件用列存聚簇布置,配合分段键和聚簇索引,范围扫描的效率会高很多。
Join小表用Hologres内部维表Join。Flink实时计算中,流和维表的关联会消耗大量资源;Hologres里如果两张事实表需要做关联,优先把小的维度表作为“复制表”或者用子查询先裁剪再关联,避免大表和大表之间的Shuffle Join。Hologres支持使用colocate表组把有相同分布键的表放在同一物理节点上,把Join变成本地关联,这一步性能提升常常是数量级的。
再来聊写入侧的调优。Hologres的实时写入走的是Flink Connector,批量发往各个Shard,默认使用行式写入,配合攒批参数可以显著降低小文件数。以下参数是我常用的一个调优组合:
-- SQL Hologres Flink Connector写入参数示例(Flink SQL) 'connector' = 'hologres', 'connection.retry.count' = '10', 'jdbc.max.retry' = '10', 'jdbc.retry.sleep' = '10s', 'mutex.retry' = '10s', 'server.timeZone' = 'Asia/Shanghai', 'write.batch.size' = '512', 'write.batch.bytes' = '10485760', 'write.batch.interval' = '1s', 'ignore.delete' = 'false', 'create.table.if.not.exists' = 'true'write.batch.size控制在单个Shard上的攒批行数,write.batch.interval是攒批窗口,write.batch.bytes是单批次字节数。这三个参数配合着调,写入吞吐会明显改善。要注意的是,攒批越大,实时可见性会降低,适合对延迟容忍度在秒级以上的场景。如果业务要求毫秒级可见,那必须关掉攒批,接受写入吞吐的下降。
Hologres对数据更新走的是Merge-On-Read机制,主键相同的记录会做就地更新。如果源数据有大量重复主键,写入和查询的负担都会加重。我踩过的一个典型坑就是:Flink作业因为状态回溯导致重复写入了海量重复主键数据,Hologres的Merge任务压力一直很高,查询性能也跟着恶化。解决方式是让Flink作业在写入前做一轮主键去重,或者利用Hologres的INSERT ... ON CONFLICT语义做幂等写入。
BI接入这一块反而是最省心的。Hologres兼容PostgreSQL协议,标准JDBC连接串是jdbc:postgresql://endpoint:port/dbname,Query BI、Tableau、帆软、DataWorks的Quick BI都能直接连。实际项目中,为了保证BI查询不拖垮实时写入链路,我会给BI报表单独建只读账号,并通过Hologres的并发控制参数限制最大连接数和Query超时时间。大促场景下,核心看板报表都走ADS层,即席查询可以放宽超时,两者互不影响。
6. 成本控制与团队转型:从离线数仓演进到实时化的一线建议
架构和技术讲得差不多了,最后聊一个很多团队不好意思摆上台面、但实际最头疼的问题:成本。Hologres的计费大头在计算资源和存储资源,实时数仓不比离线跑批,计算资源是7x24小时占用的,所以成本控制不能等到账单出来后才拍大腿。
我的成本管控经验有这几条。
第一,存储选型要匹配访问频率。Hologres支持SSD云盘和普通云盘,明细数据量大、访问频率低的分层(比如ODS历史分区),可以放在普通云盘上;高频访问的ADS、DWS层用SSD。如果预算卡得紧,ODS层甚至可以考虑在Hologres里只保留最近7天数据,更早的全量明细继续放离线数仓的冷存储里,需要回溯时再同步过去。这套“实时热数据+离线冷数据”的组合,成本能降低一半以上。
第二,计算规格要按峰值弹性扩缩容。Hologres支持计算组的弹性扩缩容,DataWorks调度也支持在夜间低峰期把实时任务切到更小的规格。我的做法是:白天业务高峰期维持大规格实例,凌晨低峰期缩容到小规格,查询类服务保留最小可用计算组保证可用性。这套弹性策略在云上体系里做起来非常简单,关键是要提前把规格变更脚本和调度任务配置好,避免人工半夜爬起来操作。
第三,实时任务的资源配额要设置上限。Flink任务在DataWorks上运行时,可以配置并行度和资源上限,防止某个实时任务出故障时把集群资源吃满,影响其他任务的写入和查询。这个“故障隔离”的思路在实时数仓里绝对要前置,不能等到线上事故再去补。
再聊聊团队转型的问题。从离线数仓转向实时数仓,最大的阻力往往不是技术,而是团队成员的操作习惯。离线工程师习惯了“写SQL、看调度、等结果”的节奏,实时化之后,需要面对Flink SQL、Checkpoint、状态、延迟监控这些新概念,心理门槛不低。
我的建议是分两步走。第一步,先把“数据同步”环节实时化,业务库的Binlog通过DataWorks实时同步直接进Hologres,让团队先练手运维Hologres表和管理数据质量,这一阶段完全不用写Flink代码。第二步,再把“计算逻辑”实时化,从最简单的滚动窗口聚合开始,用Flink SQL重写已有的离线指标,同一个指标实时和离线并行跑一段时间,对比两者结果是否一致,确认口径无误后,再逐步替换离线链路。这个过程虽然慢,但胜在稳,不会一上来就搞得团队人心惶惶。
我个人在实际操作中还有一个体会:实时数仓的上线不等于离线数仓的下线,至少在初期,“实时+离线双跑”是常态。很多团队急着把所有离线任务停掉,结果实时链路一出幺蛾子,连回退手段都没有。你真正需要的是把离线链路保留一段时间作为兜底,等实时链路的稳定性和口径验证充分了,再慢慢收缩离线任务。
最后分享一个小技巧,也是我每次新项目启动都会优先做的:把Hologres的慢Query日志和DataWorks的调度日志接入统一日志平台,设置“查询延迟超过2秒即告警”“调度任务失败即告警”的基础监控。这套东西看起来不起眼,但它在关键时刻救过我好几次——有一次深夜实时链路数据延迟量突增,就是靠着监控告警提前发现了Flink任务的反压问题,赶在业务方上班之前恢复了链路。实时数仓这个领域,稳定性永远比炫技更重要,先把监控和告警做扎实,后面所有的优化才有底气去推进。