简介:这是一套面向企业数字化转型与数据治理从业者的《大数据平台(数据中台、数据中枢、数据湖、数据要素)建设方案》演示文稿,适合架构师、数据负责人及项目规划人员用于方案汇报、框架参考与内部培训。文件围绕项目背景与目标展开,系统梳理数据中台分层架构、数据中枢的质量标准与治理能力、数据湖存储计算选型,并延伸到数据要素识别利用、技术选型实施、运维管理与持续改进等模块,目录完整,便于按章节拆解复用。资源包共1个pptx文件,约7.08MB,页面以架构图、分层要点和策略清单为主,可直接用于汇报材料或方案草稿。已有150人学习,适合需要快速搭建大数据平台建设思路、完善数据治理体系的读者参考借鉴。
1. 从一份 PPTX 建设方案说起:五个词为什么总被混着说
见过太多这样的场景:一份名为《大数据平台(数据中台、数据中枢、数据湖、数据要素)建设方案.pptx》的文档在评审会上被逐页翻过,翻到第二十页,业务方还在问“数据湖和中台到底谁管谁”。这五个词被并列写进标题,容易让人以为它们是五个可以独立立项的采购项,实际上是一条链路上的五个位置:湖解决存得下,中台解决算得准,中枢解决调得快,要素解决算得清。顺序搞反的典型后果是湖建完没人用,中台建完口径打架,中枢建完下游一改就崩。
写给正在或即将牵头这类方案的技术负责人、数据平台开发与运维,以及被拉进评审会要给出可行性判断的架构师。下面的推进方式是:先界定四个组件在链路里的位置和边界,再逐层落到会写进方案的表结构、建表语句、参数取值和上线校验脚本,中间穿插选型判断和踩坑点。
2. 数据湖、数据中台、数据中枢在大数据平台里的分层与选型
2.1 用一张分层表界定数据湖、数据中台、数据中枢的边界
很多大数据平台建设方案把四个词列成四个并列子系统,评审时才发现数据湖里存的东西和中台里的表是同一份。差别不在存什么,而在面向谁、承诺什么。数据湖面向原始与近原始明细,承诺不丢和可追溯;数据中台面向加工后的主题域数据,承诺口径一致和可复用;数据中枢面向外部调用方,承诺接口稳定和新鲜度;数据要素是加在前三者之上的一层属性,关心这份数据谁负责、值多少钱、能否被合规复用。
| 层级 | 核心职责 | 典型存储 | 时效 | 主要消费者 | 出问题的信号 |
|---|---|---|---|---|---|
| 数据湖 | 原样留存、可回溯重算 | 对象存储 + 表格式 | 小时级 | 数仓开发、算法 | 文件越存越多,无人查询 |
| 数据中台 | 主题域建模、指标统一 | 湖仓一体表 | 小时 / 分钟 | 分析师、BI | 同一指标出现三个数 |
| 数据中枢 | 服务封装、接口治理 | 结果表 + 缓存 | 秒 / 分钟 | 业务系统、应用 | 接口一改下游全崩 |
| 数据要素 | 编目、确责、计量 | 元数据表 | 天级 | 数据治理、财务 | 资产清单和实际对不上 |
写方案时这张表可以直接拿去当分工依据。判断某张表该归哪一层,问三个问题:它是不是原始数据、有没有被清洗过、有没有对外承诺接口。原始未清洗进湖,清洗后带业务含义进中台,加了稳定接口和 SLA 才进中枢。三问都答不上来的表,通常就是沉淀在 ODS 层没人管的历史包袱。
2.2 选型判断:先建数据湖还是直接上中台
判断依据不是数据量大小,而是有没有明确的分析需求、有没有能落到人头的指标负责人。数据量大但没人提需求,先建湖只能建成一个昂贵的归档盘;需求明确但数据源不稳定,先建中台则会在口径变更里反复返工。
常见的四种起点:
- 已有明确 BI 需求、指标不超过几十个、没有算法团队:可以先建轻量中台,湖只保留最薄的 ODS 落地区。
- 需要留存原始明细用于回溯重算、有外部数据接入、有模型训练诉求:先建湖,再在其上做中台,避免中台直接从业务库抽数。
- 已经有一套跑了几年的离线数仓、表结构混乱:优先做中台治理和元数据收口,湖后置,别为了“补齐链路”先上湖。
- 对时效要求高、下游是线上应用:把中枢独立出来先做服务层,指标层可以慢,接口层必须先稳。
一个容易被忽略的成本项是链路条数。每多一层,端到端延迟和故障定位成本都会增加。评审时我会建议把 ODS→DWD→DWS→ADS 的四层压缩到三层,前提是 DWD 已经承担了清洗职责,不要再单独做一个清洗中间层。
提示:方案里写“数据湖可以存一切”是危险表述,会直接导致存储预算失控。至少要在方案中写明明细保留周期和冷热分层规则。
2.3 建设方案里必须写死的三类参数
评审会上最容易被跳过、上线后最容易被追着问的,是存储、计算、时效三类参数。它们不该以“按需扩容”这种措辞出现,而要给出初始值、调整条件和责任人。
# platform-config.yaml 示例:建设方案评审时需逐项对齐 storage: lake_format: iceberg # 湖表格式,决定是否支持快照读与行级更新 file_target_size_mb: 128 # 单文件目标大小,明显低于 32MB 会触发小文件治理 retention_days: 1095 # 明细保留 3 年,超期转低频存储 cold_threshold_days: 90 # 90 天未访问转冷层 compute: engine: spark # 离线主引擎,湖表写入走它 queue_quota_cu: 400 # 队列配额,按业务域拆分 shuffle_partitions: 800 # 与数据量挂钩,默认 200 几乎必然不够 small_file_compact_cron: "0 3 * * *" # 每日凌晨合并小文件 sla: ods_to_dwd_hour: 2 # 明细层产出时限 dws_to_ads_hour: 1 # 指标层产出时限 freshness_minutes: 30 # 中枢对外承诺的数据新鲜度三类参数的逻辑不一样。存储参数决定成本曲线,文件大小和保留周期是最先被砍的两项;计算参数决定任务能不能按时跑完,shuffle_partitions给太小会 OOM、给太大会产生大量空任务;时效参数一旦对外承诺就变成合同项,宁可在方案里写宽松一点。参数说明要跟着写清调整条件,例如“单日增量超过 500GB 时把shuffle_partitions上调到 1600”,否则运维不敢动配置。
3. 数据湖落地:用 Iceberg 建一张可增量更新的明细表
3.1 Parquet 加 Iceberg 的组合为什么比直接读目录省事
Parquet 是文件格式,负责把一批行按列压好;Iceberg 是表格式,负责记录“这张表当前由哪些文件组成”。直接按dt=2024-01-01目录去读 Parquet,看似简单,但迟到数据、幂等重跑、行级撤回三件事都得人工兜。Iceberg 在数据文件之上维护 manifest 和 metadata 两层,每次写入产生一个快照,读取时按快照定位文件集合,因此能支持按快照读、按时间旅行读、行级删改和小文件合并。
方案里常被问的一句话是“我直接读 Parquet 不也能查吗”。能查,但重跑一天的分区要先把该分区文件全部删掉再写,期间查询会读到半成品。Iceberg 的写入是提交式的,读方要么看到旧快照,要么看到新快照,不会看到中间状态。对于有对账要求的场景,这一条本身就是选它的理由。
| 能力 | 裸 Parquet 目录 | Parquet + Iceberg |
|---|---|---|
| 增量读取 | 靠目录和文件时间戳猜 | 按快照 ID 精确读取 |
| 行级更新 | 不支持,只能整分区重写 | 支持 MERGE 与 DELETE |
| 并发写 | 容易相互覆盖 | 提交冲突可重试 |
| 小文件治理 | 手写脚本合并 | 内置 rewrite 动作 |
3.2 建表、写入、增量读取的具体命令
先建表。分区用days(dt)而不是dt,是为了避免每天都产生一个分区目录导致元数据膨胀;如果表按小时写入,改用hours(created_at)。
-- 在 Spark SQL 中创建一张 Iceberg 明细表 CREATE TABLE IF NOT EXISTS lake.dwd.orders ( order_id STRING COMMENT '订单号,主键', user_id STRING COMMENT '用户标识', amount DECIMAL(18,2) COMMENT '订单金额', status STRING COMMENT '订单状态', created_at TIMESTAMP COMMENT '下单时间', dt DATE COMMENT '业务日期' ) USING iceberg PARTITIONED BY (days(dt)) TBLPROPERTIES ( 'write.target-file-size-bytes' = '134217728', -- 128MB 'write.distribution-mode' = 'hash', -- 按分区键散列,避免单任务写爆 'write.metadata.delete-after-commit.enabled' = 'true' );建表语句里三个属性都影响后续运维。write.target-file-size-bytes决定单文件大小,直接关联小文件数量;write.distribution-mode设成hash后 Spark 会按分区键做 shuffle,写入并发更均匀,代价是多一次 shuffle;最后一项让旧 metadata 在提交后自动清理,不设的话元数据目录会持续膨胀。
写入与增量读取:
-- 幂等写入:先删掉当天分区再写入,避免重复跑导致数据翻倍 DELETE FROM lake.dwd.orders WHERE dt = DATE '2024-06-01'; INSERT INTO lake.dwd.orders SELECT order_id, user_id, amount, status, created_at, dt FROM lake.ods.orders_raw WHERE dt = DATE '2024-06-01';# 按快照区间读取增量数据,用于下游只处理变化部分 changed = ( spark.read.format("iceberg") .option("start-snapshot-id", "1000") .option("end-snapshot-id", "1005") .load("lake.dwd.orders") ) changed.createOrReplaceTempView("orders_delta")快照 ID 可以从lake.dwd.orders.snapshots元数据表里查。增量读取的价值在成本:下游全量扫一天分区和只读变化部分,扫描字节数可能差一个数量级。注意快照区间过大时 manifest 会膨胀,实际使用中通常把区间控制在单次调度周期内。
3.3 分区策略与文件大小:三个必调参数
分区不是越细越好。分区粒度过细会带来两个问题:一是元数据条目数量随分区数线性增长,查询计划阶段就要花掉大量时间;二是每个分区文件都很小,读取时打开大量小文件,I/O 效率反而下降。
| 参数 | 作用位置 | 建议取值 | 调错的后果 |
|---|---|---|---|
write.target-file-size-bytes | 写入端 | 128MB~256MB | 过小产生海量小文件,过大会拖长单任务 |
write.distribution-mode | 写入端 | 有分区键时用hash | 默认none会让单分区集中到一个任务 |
read.split.target-size | 读取端 | 128MB,与文件大小对齐 | 与文件大小差太多会产生大量小 split |
小文件治理要作为例行动作写进方案,而不是等出事再处理:
-- 合并指定分区的小文件,重写后文件数下降,查询扫描开销同步下降 CALL lake.system.rewrite_data_files( table => 'lake.dwd.orders', options => map( 'target-file-size-bytes', '134217728', 'min-input-files', '5' ) );min-input-files是触发阈值,只有当一个分组里的文件数超过它才重写,避免把小改动也放大成一次全分区重写。日常调度里把这条挂在凌晨低峰执行,重写完顺手跑一次expire_snapshots清理过期快照,否则存储量会因为保留大量历史快照而只涨不降。
4. 数据中台与数据中枢:指标口径统一和数据服务化
4.1 从 ODS 到 ADS 的分层命名规范与建表约束
中台最核心的产出不是表,而是统一口径。口径统一靠两件事:命名规范让每张表的归属一目了然,建表约束让不该出现的表根本建不出来。
| 层级 | 命名前缀 | 生命周期 | 是否允许对外 | 约束 |
|---|---|---|---|---|
| ODS | ods_ | 3 年 | 否 | 结构与源库一致,不做业务加工 |
| DWD | dwd_ | 3 年 | 否 | 有主键、有分区、字段有注释 |
| DWS | dws_ | 2 年 | 否 | 按主题域聚合,不含明细主键 |
| ADS | ads_ | 1 年 | 是 | 必须绑定指标定义与责任人 |
约束要落到自动化检查里。例如 DWD 层表如果没有主键注释就直接拒绝建表任务,ADS 层表如果没有登记责任人就不允许发布到中枢。规范写在文档里没人看,写在流水线的校验步骤里才会被执行。
注意:不要允许 ADS 层直接引用 ODS 层表。一旦放开,中台就退化成了一套更贵的视图集合,口径统一目标直接失效。
4.2 把指标口径写进 SQL,而不是写进文档
口径散落在文档里必然过时,写在 SQL 里才会随代码一起被评审和版本管理。做法是把指标定义收敛到 DWS 层的少数几张宽表,ADS 层只做筛选和轻聚合。
-- dws 层:每日成交指标,口径集中在这一处定义 INSERT OVERWRITE TABLE dws.daily_trade PARTITION (dt = '${bizdate}') SELECT dt, tenant_id, COUNT(DISTINCT CASE WHEN status IN ('paid','shipped') THEN order_id END) AS paid_order_cnt, SUM(CASE WHEN status IN ('paid','shipped') THEN amount ELSE 0 END) AS paid_gmv, COUNT(DISTINCT user_id) AS active_buyer_cnt FROM dwd.orders WHERE dt = '${bizdate}' GROUP BY dt, tenant_id;这段 SQL 的三个细节值得在评审时点出来。COUNT(DISTINCT ...) CASE WHEN把“支付口径”的界定收在一处,后续再加状态只需改这一行;SUM里带ELSE 0而不是直接求和,保证退款订单不会把金额带成负数;GROUP BY dt, tenant_id让多租户数据在同一张表里按租户隔离,避免每个租户建一张表。查询时对业务方只暴露ads_层,口径解释权留在数据团队。
4.3 数据中枢的服务化:一个可复用的指标查询接口
中枢的职责是把中台结果以稳定接口暴露出去,重点在参数校验、缓存和降级。下面是一个最小可用实现,重点看缓存键设计和 SQL 参数化。
import redis from fastapi import FastAPI, Query, HTTPException app = FastAPI() cache = redis.Redis(host="redis.internal", port=6379, db=0, socket_timeout=0.5) @app.get("/metrics/daily_trade") def daily_trade( dt: str = Query(..., pattern=r"^\d{4}-\d{2}-\d{2}$"), tenant: str = Query(..., max_length=32), ttl: int = Query(300, ge=60, le=3600), ): key = f"trade:{tenant}:{dt}" try: hit = cache.get(key) if hit is not None: return {"dt": dt, "tenant": tenant, "data": hit.decode(), "cache": "hit"} except redis.RedisError: pass # 缓存不可用时直接回源,不让缓存故障升级成接口故障 rows = run_sql( "SELECT paid_order_cnt, paid_gmv FROM ads.daily_trade " "WHERE dt = %(dt)s AND tenant_id = %(tenant)s", {"dt": dt, "tenant": tenant}, ) if not rows: raise HTTPException(status_code=404, detail="no data for given dt/tenant") cache.setex(key, ttl, str(rows[0])) return {"dt": dt, "tenant": tenant, "data": rows[0], "cache": "miss"}三个参数各有作用。pattern限制dt只能是日期格式,挡住注入和脏查询;ttl默认 300 秒,允许调用方在 60 到 3600 之间调整,报表类可以拉长、大屏类可以缩短;缓存键带上租户维度,避免跨租户串数据。缓存故障时直接回源而不是返回错误,这条降级逻辑在中枢里应作为默认约定,否则缓存抖动会被放大成对外接口不可用。
4.4 字段级血缘采集的落地方式
血缘采集常见两种做法:解析 SQL 文本,或者在 ETL 框架里埋点。纯解析对动态 SQL 和存储过程基本无效,纯埋点又需要所有任务都接同一套框架。务实的组合是:能解析的解析,解析不了的要求登记手写映射。
-- 血缘关系表:记录一次写入的字段级来源 CREATE TABLE IF NOT EXISTS meta.field_lineage ( job_id STRING COMMENT '任务标识', target_table STRING COMMENT '目标表', target_column STRING COMMENT '目标字段', source_table STRING COMMENT '来源表', source_column STRING COMMENT '来源字段', transform STRING COMMENT '转换表达式摘要', updated_at TIMESTAMP ) PARTITIONED BY (dt STRING);采集时机放在任务提交前,由调度平台从 SQL 计划里抽字段映射写进去。判断血缘是否可用,最简单的验收方法是抽查三张 ADS 表,人工沿着field_lineage往上追,能不能在五分钟内追到 ODS 层的原始字段。追不到就说明采集有断点,通常是某个中间任务没接进框架。
5. 数据要素视角下的进阶动作:编目、计量、校验
5.1 资产编目表与自动打标
数据要素落到工程上,第一步是把“有哪些数据、谁负责、敏感程度如何”做成可查的表,而不是散在文档里。编目表从元数据自动生成,人工只补三列:业务域、责任人、敏感等级。
INSERT OVERWRITE TABLE meta.asset_catalog PARTITION (dt = '${bizdate}') SELECT CONCAT(db_name, '.', tbl_name) AS asset_id, tbl_name AS asset_name, CASE WHEN tbl_name LIKE 'ods_%' THEN 'ODS' WHEN tbl_name LIKE 'dwd_%' THEN 'DWD' WHEN tbl_name LIKE 'dws_%' THEN 'DWS' ELSE 'ADS' END AS layer, NVL(owner, 'UNASSIGNED') AS owner, CASE WHEN tbl_name LIKE '%user%' OR tbl_name LIKE '%phone%' THEN 'SENSITIVE' ELSE 'INTERNAL' END AS sensitivity, storage_bytes, last_access_at FROM meta.table_profile WHERE dt = '${bizdate}';owner为UNASSIGNED的行数应作为治理指标每日跟踪,它直接反映有多少资产处于无人负责状态。sensitivity先按表名关键词粗筛,再人工复核,比一开始就要求逐字段打标准确率更高。
5.2 存储与计算成本的分摊口径
成本分摊要能回答“哪个业务域这个月花了多少”。口径固定为:存储按表的实际占用字节数乘单价,计算按任务消耗的 CU 时长乘单价,两者按owner汇总到业务域。分摊结果和编目表放在同一个库,方便交叉核对。需要提前约定的是共享表的处理方式,例如公共维表按调用方任务数均摊,否则每次对账都要重新扯一遍。
5.3 上线前的质量校验脚本
发布到中枢之前跑一遍校验,把问题拦在上游。下面这段检查主键唯一性和空值,结果写回质量表并在异常时阻断发布。
SELECT '${bizdate}' AS dt, COUNT(*) AS total_rows, SUM(CASE WHEN order_id IS NULL THEN 1 ELSE 0 END) AS null_pk, COUNT(*) - COUNT(DISTINCT order_id) AS dup_pk, SUM(CASE WHEN amount < 0 THEN 1 ELSE 0 END) AS negative_amount FROM ads.daily_trade WHERE dt = '${bizdate}';三个指标对应三类典型故障:null_pk非零说明上游关联丢了主键;dup_pk非零说明聚合粒度写错或多路输入没去重;negative_amount非零通常是退款口径混进了成交口径。把这三项设成流水线的阻断条件,比每次上游出问题再回查要省力得多。校验阈值不要写死在脚本里,抽到配置表中按表配置,否则每加一张表就要改一次代码。
本文还有配套的精品资源,点击获取