做大数据平台十来年,接手过的数据架构少说也有七八套。前几天帮一个团队评审他们新建的实时数仓方案,发现一个普遍问题:大家花大量时间选组件——消息队列用 Kafka 还是 Pulsar,计算引擎用 Flink 还是 Spark,OLAP 用 ClickHouse 还是 Doris——但很少有人先说清楚,这套架构到底要解决什么业务问题。数据架构的设计模式与设计原则,听起来像方法论层面的空话,实际上每一条都是从线上故障和业务事故里逼出来的经验。这篇不准备罗列组件清单,而是把大数据领域里真正经受过考验的架构模式与原则拆开讲透,适合正在设计数仓、实时计算平台,或者准备从单体数仓向湖仓一体演进的团队参考。
1. 先搞清楚数据架构在设计什么:延迟、成本与一致性的三角博弈
1.1 数据架构的职责边界
不少刚带团队的朋友问我:"数据架构是不是就是选一堆中间件搭起来?"我的回答一直是:选组件只是架构设计里最末端的一步。真正要先想明白的,是这套架构到底承担哪些职责。展开说,数据架构至少要覆盖五个环节:数据采集、数据传输、数据存储、数据计算、数据服务。
每个环节都有自己的核心矛盾。采集阶段要解决的问题是"怎么把分散在各个业务系统的数据稳定地拿过来";传输阶段要解决"数据在移动过程中不丢不重";存储阶段要解决"用什么样的格式和布局把数据放好";计算阶段要解决"用哪套引擎、以什么频率把数据加工成指标";服务阶段要解决"下游要用什么方式把数据取走"。
很多团队把架构设计做成了"组件大杂烩",就是因为把这五个环节的责任搞混了。比如为了追求实时性,把所有计算都塞进 Flink;为了查询快,又要求所有明细数据都进 ClickHouse。结果组件用了一堆,每个环节的职责边界却糊在一起,最后出了问题只能层层排查,效率极低。
我习惯把数据架构类比成城市规划:道路是传输层,仓库是存储层,工厂是计算层,商场是服务层。你不能因为商场赚钱就让整座城市全建商场,也不能因为道路修得好就让所有工厂都堆在高速路口。架构设计的本质,是给每种数据角色安排一个合理的动线和位置。
1.2 三角博弈:延迟、成本与一致性
理解了职责边界之后,第二个要面对的问题是:数据架构本质上是在做平衡。我把它总结成一个"三角博弈",这三件事几乎不可能同时做到最优:
| 目标 | 典型诉求 | 代价 |
|---|---|---|
| 低延迟 | 数据产生后秒级可见 | 计算资源密集,链路复杂,成本高 |
| 低成本 | 尽量用离线批处理、对象存储 | 数据可见性差,实时性不足 |
| 强一致 | 多路数据严格对齐、精确一次 | 需要事务、状态管理、额外网络开销 |
举个最典型的例子:业务方要求"实时大屏和离线报表的口径必须完全一致"。这句话听起来合理,但实现起来恰恰是三角博弈的核心难点。实时链路因为数据还没到齐,天然存在"数据不完整期";离线链路等数据齐了再算,结果自然不同。要强行在实时侧也做到完整期再计算,延迟就上去了;要用 OTS 那种近实时近似结果,口径就难免和离线对数对不上。
我在多个项目里验证过一件事:想清楚"哪个角可以牺牲",比盲目追求三者兼得重要得多。交易风控场景,延迟是生命线,一致性可以做到"最终一致",牺牲一点精确度换毫秒级响应;财务对账场景,一致性和准确性是底线,延迟反而可以放宽到 T+1。架构上没有银弹,只有基于业务优先级做取舍的原则。
2. Lambda、Kappa、流批一体:三种典型架构模式的选型逻辑
2.1 Lambda:双轨并行,用两套逻辑换实时性
聊大数据架构的设计模式,绕不开 Lambda。这套模式是 2011 年前后提出的,核心思路很简单:用两条独立的链路分别处理实时数据和批量数据,最终在服务层合并结果。
一条是速度层,走流处理引擎,负责处理最近几分钟甚至几秒的数据,输出低延迟的近似结果;另一条是批处理层,走 Hive/Spark 这类离线引擎,负责处理全量历史数据,输出高准确性的结果。最终查询的时候,把两条链路的结果合并起来返回给下游。
Lambda 在当年的技术背景下是合理选择,因为那时候流处理引擎的成熟度不高,精确性也差,只能靠"批处理兜底"。但真正落地过 Lambda 的人都知道它的痛点有多深:
- 两套代码要维护:同一份业务逻辑,流上写一遍,批上再写一遍,逻辑稍有不一致,口径就打架。
- 存储也有两套:速度层用 Redis 或 HBase,批处理层落在 Hive 数仓,数据冗余和一致性校验都是额外负担。
- 排查问题要跨两条链路:同一张报表数据不对,你得先判断是走了速度层还是批处理层,再分别查日志,排错成本翻倍。
我现在对团队的建议是:除非实时和离线逻辑确实差异巨大、无法统一,否则尽量别引入 Lambda。它是特定技术阶段的产物,不是值得长期坚持的架构理想态。
2.2 Kappa:用流处理统一历史与实时
Kappa 模式是 Lambda 的直接反对者:既然维护两套逻辑太痛,那我只保留一套流处理逻辑。历史数据怎么算,实时数据就怎么算。
它的核心前提是消息队列具备长时间保留数据的能力,比如 Kafka 可以把日志保留 7 天甚至更久,或者支持从对象存储里的历史日志重新回放。需要重算历史指标时,就直接从消息队列或历史日志的某个位点开始,重新跑一遍流作业,把结果写到一个新的结果表中。
Kappa 的优势非常明显:一份代码、一套引擎、统一口径。实时维度和离线维度不需要分别对口径,架构链路也短。我在一些中等数据量的业务(每天几亿到几十亿条事件)里实践过,只要把 Kafka 的 Topic 分区设计好,消息保留期拉长,Kappa 完全能撑住。
但它也有明显的限制:
- 流引擎的吞吐能力决定了全量重算的上限。如果历史数据堆积到 PB 级,纯流式重放的成本可能比离线批处理高得多。
- 复杂的 join 场景在纯流式下不好做。特别是大表 join 大表、多版本维度表关联,流处理的 状态管理 复杂度会迅速上升。
所以我的选型倾向是:数据规模在可控范围内、业务方对实时和离线口径统一要求极高的场景,优先考虑 Kappa;如果历史数据重算频次极高、规模又大,就需要引入湖仓分层来配合,而不是硬上纯 Kappa。
2.3 流批一体:当前更务实的演进方向
近几年大家谈得更多的其实是"流批一体"。它并不是简单地选 Kappa 还是 Lambda,而是在存储层和计算层同时做统一。
计算层统一用 Flink 这类引擎,同一套 SQL 既跑流模式也跑批模式,逻辑天然一致;存储层则用 Iceberg、Hudi 这类数据湖表格式,把实时写入的增量数据和离线批量写入的全量数据放在同一张表里,流批读写共用一套元数据。
我做过的项目里,比较顺手的一套组合是:
- 消息队列:Kafka,负责采集和实时传输。
- 计算引擎:Flink,负责实时 ETL、指标计算、也承担离线批任务。
- 数据湖表格式:Iceberg,承接明细数据的实时写入和离线批量合并。
- OLAP 引擎:Doris 或 ClickHouse,从 Iceberg 或直接以 Kafka 实时同步的方式供数给报表和大屏。
这套组合的好处是:实时作业和离线作业读的是同一份表数据,写的是同一个结果表,不折腾两套口径。加字段、改逻辑,只需要维护一份 Flink SQL。真正做到了用一套架构模式同时支撑实时大屏和 T+1 报表。
不过要提醒一句:流批一体对团队的要求不低,得对 Flink 的状态管理、Iceberg 的文件合并机制、以及 OLAP 引擎的实时导入链路都有比较深的理解,不然把多个复杂系统绑在一起,出了问题会很难定位。
2.4 三种模式怎么选:用业务需求倒推
我的建议是不要为了追概念选架构,而是用业务需求倒推:
- 如果业务方只要求 T+1 报表,实时指标不多,直接用数仓分层 + 离线调度即可,没必要上流批一体。
- 如果实时指标是刚需,但数据量中等、历史重算不频繁,Kappa 是性价比最高的选择。
- 如果既要求实时大屏,又要求准确的离线报表,且团队有能力驾驭多系统协作,流批一体是长期更优解。
一句话总结:架构模式是服务于业务形态的,不是拿来给简历镀金的。
3. 数据分层与存储设计:从"能存下"到"查得快"的关键决策
3.1 经典分层模型 ODS/DWD/DWS/ADS
聊完宏观架构模式,进入到具体设计层面,第一件事就是分层。数据分层是数据架构里最经典也最实用的设计原则,不同团队叫法略有出入,但骨架基本一致。
- ODS 层(操作数据层):原样接入业务系统的数据,不做加工、不丢失字段、保留历史。这一层的价值是"还原现场",任何下游口径对不上时,都能回到 ODS 层找原始数据。
- DWD 层(明细数据层):对 ODS 做清洗、去重、维度退化、统一命名,形成业务过程的事实明细。这一层讲究的是"一表一主题",比如交易明细表、订单状态流水表。
- DWS 层(汇总数据层):面向业务过程做轻度汇总,通常按主题+维度粒度设计,例如"每日商品维度交易汇总表"。这一层主要是为了减少重复计算,让下游直接查结果。
- ADS 层(应用数据层):面向具体业务需求加工的个性化数据,比如大屏指标表、数据产品接口表。
为什么几乎所有的成熟数仓都遵循这套分层?因为它的价值不是约定俗成,而是实打实的工程收益:
- 职责分离:每层只干自己的事,问题定位时知道去哪一层查。
- 复用性:DWD 层做的清洗,下游所有应用都能复用,避免每个应用各洗一遍。
- 权限管理:敏感数据在 DWD 层脱敏后,下游只能看到脱敏结果,安全边界清晰。
不过我也踩过分层过度设计的坑:某些团队把分层当成教条,每层都做全套加工,结果一条数据从 ODS 到 ADS 要跑七八个作业,延迟高、维护成本大。我现在的原则是适度分层——小微企业两三层就够,中大型团队四层为宜,每多一层就多一份延迟和运维成本。
3.2 存储选型:不同数据形态选不同引擎
很多人问我"你们数仓到底用什么存储",这个问题其实是伪命题,因为一套架构里本来就应该同时存在多种存储,各司其职。我按数据形态帮大家梳理一下:
| 数据形态 | 典型引擎 | 设计原则 |
|---|---|---|
| 实时传输管道 | Kafka / Pulsar | 分区有序、保留策略、可回放 |
| 明细数据湖 | Iceberg / Hudi / Hive | 列存、压缩、ACID、小文件治理 |
| 汇总层/OLAP | ClickHouse / Doris | 列式存储、预聚合、向量化查询 |
| 维度数据/高频查询 | Redis / HBase | 高并发点查、低延迟 |
| 原始归档 | 对象存储(S3/OSS) | 冷热分层、低成本、不可变 |
选型的核心原则,不是比哪个引擎更强,而是看访问模式。我常举一个例子:一台 OLAP 引擎再强,也不适合承接秒级高并发的点查;一个 KV 存储再快,也不适合做复杂的多维聚合分析。把数据放在它最擅长的引擎上,这是架构设计最基本的尊重。
另外一个容易被忽略的原则是尽量减少跨引擎的数据复制。早期很多团队喜欢把明细数据在 Hadoop、ClickHouse、Redis 里各放一份,看似查询快,实则一致性难题层出不穷。更推荐的做法是:明细和汇总主体只存一份,其他引擎通过同步/物化的方式按需分发,并明确主从关系。
3.3 分区、分桶与排序:查询性能的设计前置
存储选完,还得在表设计上花功夫。分区、分桶、排序这三件事直接影响查询性能,而且是"设计期决定、运行期难改"的。
分区策略上,绝大多数场景按时间分区,比如 dt=2024-06-01。时间分区的好处是自然贴合业务查询(查某一天、某一个月的报表),也方便生命周期管理,直接删分区就是删数据。但要注意,如果业务经常按用户维度过滤,而数据量又很大,最好在用户维度上再设计分桶或者直接使用 OLAP 引擎的明细表索引。
排序设计容易被忽视。以 ClickHouse 为例,表引擎的排序列决定了索引的稀疏结构,排序列选错,查询性能可能差一个数量级。我的经验是:把高频过滤字段和聚合维度尽量往前放,比如按时间、用户、渠道排序,而不是按一个几乎不用来过滤的字段排序。
还有一个高频问题:小文件。流式写入数仓时,如果不做 compaction,几分钟就会产生一个小文件。小文件多了,NameNode 压力大、查询扫描效率低。应对方案是在表格式层面开启自动合并(Iceberg 的 compaction、Hudi 的 clustering),或者控制写入并行度,让每个文件尽量到达合适大小。
3.4 数据模型设计:维度建模与宽表
最后说数据模型。传统数仓里经典的星型模型、雪花模型,在大数据体系里依然适用,但在实操中我越来越多地偏向宽表化。
原因很直接:宽表把多个维度的字段冗余到一张表里,查询时不需要多表 join,ETL 时也只需要算一次。对于 OLAP 引擎而言,宽表扫描效率往往优于频繁 join,尤其在海量明细场景下,减少 join 就是减少 shuffle 和状态管理开销。
但宽表也不是越宽越好。字段过多会导致表结构臃肿、Schema 演化困难。我的设计原则是:一个主题域一张宽表,字段控制在业务指标真正需要的范围内,而不是把上游所有字段全部堆进来。冗余时要问自己一个问题:这个字段真的会被下游频繁使用吗?如果不是,宁可留在明细层,需要时再关联。
4. 实时数据链路里的幂等、一致性与容错设计:细节决定成败
4.1 端到端延迟的构成
架构模式选好了,表也建好了,接下来要处理的是实时链路里的工程细节。很多人以为 Flink 算得快就是实时性好,其实端到端延迟是一个链条上的累积:
数据从业务系统产生,经过日志采集或 binlog 采集进 Kafka,Flink 消费后做清洗和 join,再写入下游 OLAP,最后在大屏上展示。这五个环节的延迟加起来,才是业务真正感知到的延迟。
我在真实项目里见过不少"伪实时"案例:Flink 作业的 processing time 只有几百毫秒,但上游采集端每 10 分钟才 Flush 一批数据,下游 OLAP 写入端又是攒批写入,结果大屏数据永远滞后 10 分钟以上。排查下来才发现瓶颈根本不在计算引擎。
优化端到端延迟的正确思路,是逐步测量每个环节的 p95、p99 延迟,找出最长的一环。通常最容易出问题的是采集端的攒批策略、Kafka 分区数不够导致消费并行度上不去、以及 OLAP 写入端为了吞吐设置过大的攒批阈值。这三个位置调优好,端到端延迟会明显改善。
4.2 幂等写入与数据去重
实时链路中,数据重复几乎是必然事件。上游重试发送、Flink 重启后从 checkpoint 恢复重新处理、下游写入遇到网络超时重试,每个环节都可能产生重复数据。所以架构设计必须默认"会有重复",并设计去重和幂等方案。
去重设计的第一原则是找到业务自然主键。例如订单事件可以用 order_id,交易流水可以用流水号。在写入结果表时,以主键做 upsert:存在则更新,不存在则插入。以 Flink + Iceberg 为例,Iceberg 本身支持主键 upsert,写入时指定主键列即可;Flink + Kafka + OLAP 的场景,则通常依赖 OLAP 引擎的 Unique 模型或者通过状态去重后写入。
但要注意一个坑:有些业务根本没有自然主键,比如用户行为日志,多条记录可能完全一样。这时候需要在采集端追加一个全局唯一标识字段,比如 UUID 或者"日志时间 + 随机数"的组合。否则下游无论如何做去重都没有抓手。
另外,状态去重时要关注状态大小。用 Flink 的 ValueState 或者 RocksDB 做去重,状态会随基数增长,超过内存后性能会退化。我见过一个项目,为了去重把几十亿 key 都放进了状态后端,结果作业频繁 GC、checkpoint 超时。更合理的方案是:基于时间窗口做去重,老数据让上游重放或依赖下游合并,不要试图在流状态里保存所有历史 key。
4.3 Exactly-Once 的正确理解
很多团队一提实时计算就要求"精确一次语义",但Exactly-Once 不是单靠 Flink 就能完成的,而是端到端的整体设计。
Flink 的 checkpoint 机制能保证的是"计算引擎内部的精确一次":算子状态和 Kafka offset 被原子地提交,重启后不会重复计算。但计算完写入下游时,如果下游不支持幂等,依然可能出现重复写入。所以完整的端到端精确一次,必须满足两个条件:
- 上游数据源支持位点重置(Kafka offset 可提交、可回退)。
- 下游存储支持幂等写入(Iceberg/Hudi 的主键能力、或者 OLAP 的 Unique 模型)。
我常用的判断标准是:_如果下游存储只支持 append 写入,那么所谓的 Exactly-Once 只是计算引擎层面的,最终效果仍然是 At-Least-Once。所以做实时链路之初,就要把下游存储的语义确认清楚,而不是等上线后才发现数据多算了一倍。
4.4 乱序数据处理:watermark 的工程经验
实时流式数据天然存在乱序:用户点击、支付事件在各个节点上经过不同的网络路径,到达 Kafka 的先后顺序可能和事件发生顺序不一致。如果不处理乱序,统计结果就会出现"某分钟成交额少算、下一分钟又突然多算"的现象。
Flink 解决乱序的标准工具是 watermark + allowedLateness。watermark 是"数据完整性的进度信号",假设我们允许事件最多迟到 10 秒,那么 10 秒前的数据视为可输出;如果某条 10 秒前的事件到得更晚,就会被丢弃或进入侧输出流。
这里分享一个调参经验:watermark 的延迟时间不要拍脑袋定,而是要根据线上数据延迟分布的 p95 或 p99 来设置。设置太短,乱序数据被大量丢弃,指标偏差大;设置太长,实时性受损。我在订单分析场景里一般取延迟分布 p99 的值,再稍加一点余量,既保证绝大部分数据能归位,又不至于让指标滞后太多。
还要留意一个常见陷阱:watermark 是和 source 分区绑定的。如果 Kafka 某个分区始终没有数据,会造成 watermark 不推进,整个聚合被卡住。应对方案是配合空闲分区超时机制(Idle Source),让长时间无数据的分区不阻塞整体水位推进。
5. 元数据与数据治理:数据架构的骨架与"数据地图"
5.1 元数据不止是"建表语句"
很多团队建好数仓、接通实时链路后,以为架构设计就完成了。直到有一天业务方拿着新需求问"我们的用户活跃数据到底在哪个表里、口径是什么、能不能直接对接",整个团队才发现:没有元数据管理,数仓就是个黑盒。
元数据至少分三层:
- 技术元数据:表结构、分区信息、存储路径、字段类型、依赖关系。
- 业务元数据:指标定义、口径说明、负责人、业务归属。
- 管理元数据:数据质量规则、数据生命周期、访问权限、脱敏策略。
元数据管理是数据架构的骨架。没有这套骨架,数据鲜活性再高也无法被有效率地使用。在我参与的项目里,元数据系统一般承担三个核心职责:自动采集表结构,沉淀指标口径,支持检索和查看血缘。
实操层面,如果不想引入重量级元数据平台,也可以先从"元数据表 + 自动同步脚本"开始:把 Hive/Iceberg 的表元数据定时同步到 MySQL 中,再用开源的 DataHub / Amundsen 或自研一个简单的检索页面。这些工作早做比晚做强得多,后期表数量超过几百张后再补,成本会成倍上升。
5.2 数据血缘:排查问题的关键基础设施
元数据里最有价值也最不好做的,是数据血缘。
血缘就是一张"数据流向地图":ODS 层的表 A 经过哪个作业生成了 DWD 层的表 B,DWD 层的表 B 又供给了 DWS 层的表 C。没有血缘,每次数据异常排查都像在黑暗中摸箱子。
血缘的生成方式有两条路径:解析 SQL 自动生成和从调度系统采集任务依赖间接推断。前者准确度高,但对复杂 SQL 的自研解析有门槛;后者容易实现,但粒度偏粗,只能看到任务级依赖,看不清字段级流向。
我自己的实践经验是:先靠调度平台已有的任务依赖把表级血缘跑通,再逐步补充字段级血缘。日常排查中最常用的场景是这样的:业务方说报表里的"昨日成交额"和财务口径对不上,有了血缘,就能从 ADS 层一步步回溯到 DWD 层,找到是哪一步加工逻辑引入了偏差,而不是靠记忆去翻十几个脚本。
5.3 数据质量监控与告警策略
架构做得再好,数据质量出问题,一切归零。质量监控应该在架构设计阶段就留出位置,而不是等出了问题再补。
我按优先级整理几个核心监控项:
- 完整性监控:每个任务周期内,表的数据量是否达到预期。写一个基线任务,统计 ODS 表当日记录数和前 7 天均值做对比,波动超过 20% 就要告警。
- 空值/异常值监控:核心字段(订单金额、用户 ID)空值率一旦超过阈值,大概率是上游采集有问题。
- 及时性监控:实时链路中,数据写入落库的时间和事件时间的时间差是否稳定,如果延迟突然上窜,需要马上关注上游攒批或下游写入瓶颈。
- 口径一致性校验:定期用同一指标在实时结果和离线结果之间做对比,差异超过阈值时触发检查。这个通常是数仓团队最头疼但也最重要的监控。
告警策略上切忌"全盘告警"。我刚带团队时把每个检查都设成告警,结果一天几百条通知,应急响应成员很快就麻木了。合理做法是分级:red 级告警(数据不可用、严重延迟)实时通知到人,yellow 级告警(数据波动、小范围空值)只记录到日报,第二天统一处理。
5.4 权限与安全:架构的底线
数据架构设计里,权限和安全经常被排在最后,直到真的出事才被想起。在数据量越来越大、敏感字段越来越多的今天,权限设计应该前置。
我推荐的基本框架是:
- 库表级权限:通过统一的权限中心控制,不同角色只能访问自己业务域的表。
- 字段级脱敏:手机号、身份证、支付账号等敏感字段在 ODS 层就要制定脱敏策略,DWD 层及以下使用脱敏后的数据。
- 操作审计:谁在什么时间查询了哪张表、导出了多少行数据,要有完整的日志留痕。
权限体系的落地不复杂,复杂的是和业务流程结合:例如"运营人员能看汇总值但不能看明细","风控人员能看客户明细但不能导出"等,这类规则需要在架构设计和权限模型里提前支持,否则后期强行打补丁会非常痛苦。
6. 实操复盘:一个实时数仓数据架构的完整设计过程
6.1 需求梳理:先定义问题
理论讲了不少,我拿最近做的一个项目完整复盘一遍,把上面说的模式、原则、细节串起来。背景是一家互联网电商公司,业务方提了几个需求:
- 实时大屏:要看到今日实时 GMV、订单量、各省份销售分布,刷新延迟要求 5 秒以内。
- 实时预警:当某个爆品库存低于阈值时,运营要立刻收到通知。
- 离线报表:财务、商品、供应链各域仍然保留 T+1 的常规日报和周报。
- 口径要求:实时指标和离线指标在 T+1 后必须能对上数。
初步估算数据量:核心业务事件每天约 5 亿条,峰值 TPS 约 3 万,明细数据日均新增 2TB 左右。这是一个典型的"既要实时又要离线、还要口径统一"的场景。
6.2 架构选型:组件与模式
基于需求,我直接排除了纯 Lambda(不想要两套逻辑),也不选择纯 Kappa(历史重算规模太大,纯流式重放成本高)。最终确定的方案是流批一体为主,分层存储配合:
- 采集层:业务日志走 Kafka,数据库变更走 Flink CDC 同步到 Kafka。
- 存储层:明细数据统一落到 Iceberg 表,按天分区;同时以 Kafka 实时链路将明细同步到 Doris。
- 计算层:Flink 实时计算指标,产出实时汇总结果;离线部分用 Flink Batch 模式跑 DWD/DWS 层加工,共用同一份 SQL 逻辑。
- 服务层:Doris 面向实时大屏和即席查询;Hive/Iceberg 面向离线报表和数据科学团队。
- 元数据与调度:用 DataHub 采集表血缘,用 DolphinScheduler 编排离线任务。
这套架构里,Flink 的实时作业和离线作业读的是同一张 Iceberg 表,写的是同一个 DWS 层结果表,口径天然一致。实时链路则通过 Kafka → Doris 的方式保障秒级大屏需求。
6.3 关键设计决策与理由
几个关键决策我在评审会上一一说明过,这里也列给大家参考:
为什么明细层用 Iceberg 而不用纯 Hive?因为 Iceberg 支持 ACID、行级 upsert、隐藏分区和高效的 compaction,能让实时写入和离线批量写入在同一张表上安全并存。Hive 表做实时写入,小文件问题和并发冲突会比较难处理。
为什么实时结果直接写 Doris 而不是先落 Iceberg 再同步?因为大屏对延迟的要求是 5 秒内,从 Kafka 进 Doris 的实时导入链路最短、延迟最低。离线报表则从 Iceberg 算,最后对不上数时用 DWS 层的同一份结果做校验口径。简单说:实时走短链路,离线走长链路,两条链路在结果表层面收敛。
实时和离线的一致性如何保证?实时作业和离线作业是用同一份 Flink SQL 模板生成的两个 Job,最大粒度、过滤条件、聚合维度完全一致。校验方式是每日凌晨用离线 DWS 结果覆盖一次 Doris 对应分区的数据,彻底消除实时累积误差。这就是一种"离线校正实时"的兜底策略。
6.4 落地效果与后续演进
这套架构上线后,大屏指标延迟稳定在 3 秒左右,p95 为 4.5 秒,满足需求;T+1 口径比对,连续一个月误差低于万分之一,远超业务预期。
当然也存在几个后续演进点:一是 Doris 里的明细数据保存周期需要根据成本重新评估;二是 Iceberg 的 compaction 在小文件高峰时会给集群带来压力,后续考虑把 compaction 调度错峰执行;三是数据湖权限体系还没完全接入,目前只能在应用层做控制,后面要推进统一 Ranger 认证。
7. 我在多个项目里踩过的架构坑:设计文档里不会写的教训
7.1 过度设计:性能瓶颈问题还没出现,先把最简单的方案做扎实
我早期做架构有个毛病:总想把最先进、最复杂的方案一次到位。有一次项目日数据量才几千万,我就把流批一体、湖仓分离、数据湖格式全部上了。结果运维复杂度远超团队承受能力,一个简单需求要动四五个系统。
后来我给自己定了条规矩:数据量和业务复杂度没有到那个量级,就不要用那个量级的架构。先用最简单的方案把链路跑顺——比如离线数仓 + 定时调度,等确有实时需求且数据量上来了,再逐步演进。架构是演进而来的,不是一步到位的。
7.2 小文件与大表:存储设计偷的懒,会在半年后加倍还回来
另一个我踩得比较深的坑是忽视了小文件治理。当时一个流式作业每 5 分钟写一次结果表,每次写几百个小文件,一个月后表内文件数超过十几万,查询扫描效率直线下降,甚至影响到了集群整体性能。
后来在表格式层面开了自动 compaction,并人为控制写入并行度——数据量没到那么大的时候,并行度不是越高越好。保持每个 Spark 或 Flink 写入任务生成的文件尽量接近目标大小,比事后一遍遍整理小文件高效得多。
7.3 忽略数据回溯能力:架构必须为"算错了"留退路
数据架构里最容易被忽略的,是回溯能力。很多团队把链路搭好就不管了,直到某天发现口径错误,需要重算过去 30 天的数据,才意识到根本没有设计重算机制。
要解决这个问题,核心是数据保留策略和表格式支持。Iceberg/Hudi 这类支持时间旅行的表格式是很好的基础,可以基于某个快照重新派生数据,而不必完全重放所有原始日志。另一个是所有加工逻辑必须参数化,至少支持指定日期范围重跑,别把逻辑写死在调度脚本里。
7.4 预留 Schema 演化空间:加一个字段引发的线上事故
有过一次特别深刻的教训:业务方临时要求在一个大表里加字段,但表结构当时是强约束,改动涉及下游几十个任务的重启,最后花了整整两天才全部对齐。从那以后,我在表设计时都要求预留 schema 演化空间:要么采用支持 schema evolution 的表格式(Iceberg/Hudi 是天然的),要么约定加字段时只允许后置,不允许改变已有字段的语义,并且建立字段评审机制。
7.5 团队能力与架构复杂度要匹配
最后一个坑不在技术上,在团队组织上。再好的架构,如果团队没人能驾驭,就是灾难。流批一体里涉及的 Flink 状态管理、Watermark 调优、Iceberg compact 策略,都需要比较专门的技能。上线前要评估团队是否具备这些能力储备,如果还比较薄弱,宁可先从相对简单的架构开始,再安排培训和逐步复杂化。
我自己现在的态度是:架构方案的第一评审人不是技术负责人,而是未来半年负责运维这套系统的一线同学。他们说"看不懂、hold 不住"的地方,就该简化,或者先补齐能力再上。毕竟数据架构最终的衡量标准不是方案多先进,而是线上系统是否稳定,业务问题是否被高效解决。