尧图精选

大数据架构设计模式与原则:实时数仓与流批一体选型实战

🕒 发布时间:2026/10/2 14:06:49 📁 来源:尧图网络
做大数据平台十来年接手过的数据架构少说也有七八套。前几天帮一个团队评审他们新建的实时数仓方案发现一个普遍问题大家花大量时间选组件——消息队列用 Kafka 还是 Pulsar计算引擎用 Flink 还是 SparkOLAP 用 ClickHouse 还是 Doris——但很少有人先说清楚这套架构到底要解决什么业务问题。数据架构的设计模式与设计原则听起来像方法论层面的空话实际上每一条都是从线上故障和业务事故里逼出来的经验。这篇不准备罗列组件清单而是把大数据领域里真正经受过考验的架构模式与原则拆开讲透适合正在设计数仓、实时计算平台或者准备从单体数仓向湖仓一体演进的团队参考。1. 先搞清楚数据架构在设计什么延迟、成本与一致性的三角博弈1.1 数据架构的职责边界不少刚带团队的朋友问我数据架构是不是就是选一堆中间件搭起来我的回答一直是选组件只是架构设计里最末端的一步。真正要先想明白的是这套架构到底承担哪些职责。展开说数据架构至少要覆盖五个环节数据采集、数据传输、数据存储、数据计算、数据服务。每个环节都有自己的核心矛盾。采集阶段要解决的问题是怎么把分散在各个业务系统的数据稳定地拿过来传输阶段要解决数据在移动过程中不丢不重存储阶段要解决用什么样的格式和布局把数据放好计算阶段要解决用哪套引擎、以什么频率把数据加工成指标服务阶段要解决下游要用什么方式把数据取走。很多团队把架构设计做成了组件大杂烩就是因为把这五个环节的责任搞混了。比如为了追求实时性把所有计算都塞进 Flink为了查询快又要求所有明细数据都进 ClickHouse。结果组件用了一堆每个环节的职责边界却糊在一起最后出了问题只能层层排查效率极低。我习惯把数据架构类比成城市规划道路是传输层仓库是存储层工厂是计算层商场是服务层。你不能因为商场赚钱就让整座城市全建商场也不能因为道路修得好就让所有工厂都堆在高速路口。架构设计的本质是给每种数据角色安排一个合理的动线和位置。1.2 三角博弈延迟、成本与一致性理解了职责边界之后第二个要面对的问题是数据架构本质上是在做平衡。我把它总结成一个三角博弈这三件事几乎不可能同时做到最优目标典型诉求代价低延迟数据产生后秒级可见计算资源密集链路复杂成本高低成本尽量用离线批处理、对象存储数据可见性差实时性不足强一致多路数据严格对齐、精确一次需要事务、状态管理、额外网络开销举个最典型的例子业务方要求实时大屏和离线报表的口径必须完全一致。这句话听起来合理但实现起来恰恰是三角博弈的核心难点。实时链路因为数据还没到齐天然存在数据不完整期离线链路等数据齐了再算结果自然不同。要强行在实时侧也做到完整期再计算延迟就上去了要用 OTS 那种近实时近似结果口径就难免和离线对数对不上。我在多个项目里验证过一件事想清楚哪个角可以牺牲比盲目追求三者兼得重要得多。交易风控场景延迟是生命线一致性可以做到最终一致牺牲一点精确度换毫秒级响应财务对账场景一致性和准确性是底线延迟反而可以放宽到 T1。架构上没有银弹只有基于业务优先级做取舍的原则。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。真正做到了用一套架构模式同时支撑实时大屏和 T1 报表。不过要提醒一句流批一体对团队的要求不低得对 Flink 的状态管理、Iceberg 的文件合并机制、以及 OLAP 引擎的实时导入链路都有比较深的理解不然把多个复杂系统绑在一起出了问题会很难定位。2.4 三种模式怎么选用业务需求倒推我的建议是不要为了追概念选架构而是用业务需求倒推如果业务方只要求 T1 报表实时指标不多直接用数仓分层 离线调度即可没必要上流批一体。如果实时指标是刚需但数据量中等、历史重算不频繁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、小文件治理汇总层/OLAPClickHouse / Doris列式存储、预聚合、向量化查询维度数据/高频查询Redis / HBase高并发点查、低延迟原始归档对象存储S3/OSS冷热分层、低成本、不可变选型的核心原则不是比哪个引擎更强而是看访问模式。我常举一个例子一台 OLAP 引擎再强也不适合承接秒级高并发的点查一个 KV 存储再快也不适合做复杂的多维聚合分析。把数据放在它最擅长的引擎上这是架构设计最基本的尊重。另外一个容易被忽略的原则是尽量减少跨引擎的数据复制。早期很多团队喜欢把明细数据在 Hadoop、ClickHouse、Redis 里各放一份看似查询快实则一致性难题层出不穷。更推荐的做法是明细和汇总主体只存一份其他引擎通过同步/物化的方式按需分发并明确主从关系。3.3 分区、分桶与排序查询性能的设计前置存储选完还得在表设计上花功夫。分区、分桶、排序这三件事直接影响查询性能而且是设计期决定、运行期难改的。分区策略上绝大多数场景按时间分区比如 dt2024-06-01。时间分区的好处是自然贴合业务查询查某一天、某一个月的报表也方便生命周期管理直接删分区就是删数据。但要注意如果业务经常按用户维度过滤而数据量又很大最好在用户维度上再设计分桶或者直接使用 OLAP 引擎的明细表索引。排序设计容易被忽视。以 ClickHouse 为例表引擎的排序列决定了索引的稀疏结构排序列选错查询性能可能差一个数量级。我的经验是把高频过滤字段和聚合维度尽量往前放比如按时间、用户、渠道排序而不是按一个几乎不用来过滤的字段排序。还有一个高频问题小文件。流式写入数仓时如果不做 compaction几分钟就会产生一个小文件。小文件多了NameNode 压力大、查询扫描效率低。应对方案是在表格式层面开启自动合并Iceberg 的 compaction、Hudi 的 clustering或者控制写入并行度让每个文件尽量到达合适大小。3.4 数据模型设计维度建模与宽表最后说数据模型。传统数仓里经典的星型模型、雪花模型在大数据体系里依然适用但在实操中我越来越多地偏向宽表化。原因很直接宽表把多个维度的字段冗余到一张表里查询时不需要多表 joinETL 时也只需要算一次。对于 OLAP 引擎而言宽表扫描效率往往优于频繁 join尤其在海量明细场景下减少 join 就是减少 shuffle 和状态管理开销。但宽表也不是越宽越好。字段过多会导致表结构臃肿、Schema 演化困难。我的设计原则是一个主题域一张宽表字段控制在业务指标真正需要的范围内而不是把上游所有字段全部堆进来。冗余时要问自己一个问题这个字段真的会被下游频繁使用吗如果不是宁可留在明细层需要时再关联。4. 实时数据链路里的幂等、一致性与容错设计细节决定成败4.1 端到端延迟的构成架构模式选好了表也建好了接下来要处理的是实时链路里的工程细节。很多人以为 Flink 算得快就是实时性好其实端到端延迟是一个链条上的累积数据从业务系统产生经过日志采集或 binlog 采集进 KafkaFlink 消费后做清洗和 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 层的表 BDWD 层的表 B 又供给了 DWS 层的表 C。没有血缘每次数据异常排查都像在黑暗中摸箱子。血缘的生成方式有两条路径解析 SQL 自动生成和从调度系统采集任务依赖间接推断。前者准确度高但对复杂 SQL 的自研解析有门槛后者容易实现但粒度偏粗只能看到任务级依赖看不清字段级流向。我自己的实践经验是先靠调度平台已有的任务依赖把表级血缘跑通再逐步补充字段级血缘。日常排查中最常用的场景是这样的业务方说报表里的昨日成交额和财务口径对不上有了血缘就能从 ADS 层一步步回溯到 DWD 层找到是哪一步加工逻辑引入了偏差而不是靠记忆去翻十几个脚本。5.3 数据质量监控与告警策略架构做得再好数据质量出问题一切归零。质量监控应该在架构设计阶段就留出位置而不是等出了问题再补。我按优先级整理几个核心监控项完整性监控每个任务周期内表的数据量是否达到预期。写一个基线任务统计 ODS 表当日记录数和前 7 天均值做对比波动超过 20% 就要告警。空值/异常值监控核心字段订单金额、用户 ID空值率一旦超过阈值大概率是上游采集有问题。及时性监控实时链路中数据写入落库的时间和事件时间的时间差是否稳定如果延迟突然上窜需要马上关注上游攒批或下游写入瓶颈。口径一致性校验定期用同一指标在实时结果和离线结果之间做对比差异超过阈值时触发检查。这个通常是数仓团队最头疼但也最重要的监控。告警策略上切忌全盘告警。我刚带团队时把每个检查都设成告警结果一天几百条通知应急响应成员很快就麻木了。合理做法是分级red 级告警数据不可用、严重延迟实时通知到人yellow 级告警数据波动、小范围空值只记录到日报第二天统一处理。5.4 权限与安全架构的底线数据架构设计里权限和安全经常被排在最后直到真的出事才被想起。在数据量越来越大、敏感字段越来越多的今天权限设计应该前置。我推荐的基本框架是库表级权限通过统一的权限中心控制不同角色只能访问自己业务域的表。字段级脱敏手机号、身份证、支付账号等敏感字段在 ODS 层就要制定脱敏策略DWD 层及以下使用脱敏后的数据。操作审计谁在什么时间查询了哪张表、导出了多少行数据要有完整的日志留痕。权限体系的落地不复杂复杂的是和业务流程结合例如运营人员能看汇总值但不能看明细风控人员能看客户明细但不能导出等这类规则需要在架构设计和权限模型里提前支持否则后期强行打补丁会非常痛苦。6. 实操复盘一个实时数仓数据架构的完整设计过程6.1 需求梳理先定义问题理论讲了不少我拿最近做的一个项目完整复盘一遍把上面说的模式、原则、细节串起来。背景是一家互联网电商公司业务方提了几个需求实时大屏要看到今日实时 GMV、订单量、各省份销售分布刷新延迟要求 5 秒以内。实时预警当某个爆品库存低于阈值时运营要立刻收到通知。离线报表财务、商品、供应链各域仍然保留 T1 的常规日报和周报。口径要求实时指标和离线指标在 T1 后必须能对上数。初步估算数据量核心业务事件每天约 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 秒满足需求T1 口径比对连续一个月误差低于万分之一远超业务预期。当然也存在几个后续演进点一是 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 不住的地方就该简化或者先补齐能力再上。毕竟数据架构最终的衡量标准不是方案多先进而是线上系统是否稳定业务问题是否被高效解决。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →