Lesson 06 · 落地
把 PG 冷数据分层到 Parquet
一条真跑通的管线:51 MB 的 PG 表 → 3.16 MB 的分区 Parquet,行数与校验和一致。
前五课的知识现在要变成一个方案。这一课给出 Pivot 的冷数据分层骨架、可直接改用的导出命令、以及三个"什么时候该升级到 Iceberg"的判据。本课所有命令都在一个临时 PostgreSQL 16 实例上真跑过,不是伪代码。
分层架构
写入 查询
│ │
┌─────▼──────────────┐ ┌────▼─────────────────────────────┐
│ 热层 · PostgreSQL │◀───│ 最近 N 天:设备详情、实时看板、告警 │
│ 最近 30–90 天 │ │ 点查、更新、事务,全部在这层 │
│ 行存 + 索引 + 事务 │ └──────────────────────────────────┘
└─────┬──────────────┘
│ 每日归档作业(幂等,可重跑)
│ ① 导出 → ② 校验 → ③ 删源
┌─────▼──────────────┐ ┌──────────────────────────────────┐
│ 冷层 · Parquet on │◀───│ 历史分析:季度能耗曲线、同比、报表 │
│ S3 / OSS │ │ DuckDB 直接查,或下载到本地算 │
│ day 分区 + 内部排序 │ └──────────────────────────────────┘
└────────────────────┘
三条边界要划清楚:
- 热层的职责不变。所有需要事务、更新、主键点查的东西留在 PG。冷层不接管任何在线读路径。
- 冷层是只读的。要修历史数据,就整个分区重算重写——这是可接受的,因为 IoT 读数本来就不该被改。
- 分界线是时间,不是数据量。一条明确的规则("90 天前的进冷层")比"表太大了就归档"好维护得多。
一条命令完成导出
DuckDB 的 postgres 扩展可以直接把 PG 当成一个库挂上来,然后用一条 COPY 完成"查询 + 排序 + 分区 + 压缩 + 落盘":
INSTALL postgres; LOAD postgres; -- 只读挂载,避免任何误写风险 ATTACH 'host=… port=5432 dbname=pivot user=… password=…' AS pg (TYPE postgres, READ_ONLY); COPY ( SELECT ts, tenant_id, building_code, device_id, metric, value, quality, strftime(ts AT TIME ZONE 'UTC', '%Y-%m-%d') AS day -- 显式 UTC,见第 5 课时区坑 FROM pg.public.iot_reading WHERE ts >= TIMESTAMPTZ '2026-01-01 00:00:00+00' AND ts < TIMESTAMPTZ '2026-01-16 00:00:00+00' -- 左闭右开,别用 BETWEEN ORDER BY device_id, ts -- 第 4 课:排序键决定跳过 ) TO 's3://pivot-archive/iot' (FORMAT parquet, PARTITION_BY (day), COMPRESSION zstd, ROW_GROUP_SIZE 200000);
行数 432,000 与 sum(value) 校验和两端完全一致。体积比 16 : 1——注意 PG 侧的 51 MB 包含索引和每行 23 字节的元组头,这部分开销在冷层彻底消失。
归档作业的骨架
把上面这条命令包成每日作业时,三件事决定它能不能安全地跑在生产上:
import duckdb, sys DAY = sys.argv[1] # 例:2026-01-08,由调度器传入 DST = f"s3://pivot-archive/iot/day={DAY}" con = duckdb.connect() con.sql("INSTALL postgres; LOAD postgres; INSTALL httpfs; LOAD httpfs;") con.sql("ATTACH '…' AS pg (TYPE postgres, READ_ONLY)") # ① 幂等:先写到临时前缀,成功后再改名/移动到正式路径 tmp = DST + ".tmp" con.sql(f""" COPY (SELECT … FROM pg.public.iot_reading WHERE ts >= TIMESTAMPTZ '{DAY} 00:00:00+00' AND ts < TIMESTAMPTZ '{DAY} 00:00:00+00' + INTERVAL 1 DAY ORDER BY device_id, ts) TO '{tmp}' (FORMAT parquet, COMPRESSION zstd, ROW_GROUP_SIZE 200000)""") # ② 校验:行数 + 关键列校验和,两端必须相等 src = con.sql(f"SELECT count(*), sum(value) FROM pg.public.iot_reading WHERE …").fetchone() dst = con.sql(f"SELECT count(*), sum(value) FROM read_parquet('{tmp}/**/*.parquet')").fetchone() assert src == dst, f"校验失败 src={src} dst={dst}" # ③ 只有校验通过才提交,并且删源要单独一步、可延迟执行 # rename(tmp → DST);删 PG 数据建议延后 7 天,给回滚留窗口
- 幂等:作业重跑必须得到同样结果,不能追加出重复数据。手段是"临时路径 + 整体提交",不是"先删后写"。
- 校验:行数 + 数值列求和,两端比对。这是最便宜的正确性保证,别省。
- 删源延后:导出成功不等于可以立刻删 PG。留 7 天窗口,代价只是一点磁盘。
查询侧:冷数据怎么被读到
INSTALL httpfs; LOAD httpfs; CREATE SECRET (TYPE s3, KEY_ID '…', SECRET '…', REGION '…', ENDPOINT '…'); SELECT day, device_id, round(avg(value), 2) AS avg_kw FROM read_parquet('s3://pivot-archive/iot/**/*.parquet', hive_partitioning = true) WHERE day BETWEEN '2026-01-01' AND '2026-03-31' AND building_code = 'BLD-001' GROUP BY 1, 2 ORDER BY 1;
冷热合并查询有两条路,按团队情况选:
| 做法 | 适合 | 代价 |
|---|---|---|
| 应用层按时间路由 查询时间窗 < 90 天走 PG,否则走 DuckDB |
大多数团队的起点 | 要在服务里写一段路由逻辑,跨界查询要合并两边结果 |
| DuckDB 同时挂 PG 和 S3 用一条 UNION ALL 视图统一冷热 |
报表 / BI 这类离线场景 | DuckDB 成为一个需要运维的组件;不适合承担在线请求 |
什么时候该上 Iceberg
裸 Parquet 目录(本课的方案)够用到什么时候?三个信号,出现任意两个就该认真评估表格式了:
- 并发写冲突:不止一个作业在写同一批分区,或者 compaction 与写入撞车(第 5 课那个非原子的 compaction)。
- 需要行级更新/删除:合规要求删除某租户的历史数据、或要修正一批错误读数。裸目录只能整分区重写。
- schema 演进超出"加可空列":要改列名、改类型、或需要"回到上周的数据状态"(时间旅行)。
反过来说:如果只是"一个作业每天往里写一次、只加列、不删改",裸 Parquet 目录完全够用,上 Iceberg 只是白白增加一套元数据服务要维护。别为了名字上湖仓。
2026 年值得知道的两个新东西
Variant:给"每个设备字段都不一样"的 payload 用
Pivot 一定会遇到这种数据:不同厂商的设备上报的 JSON 结构完全不同,字段还随固件版本变。传统做法是整个 JSON 存成一个字符串列——查询时每次都要解析全文,也没法做统计信息。
Parquet 的 Variant 类型(2025 年 8 月引入规范,2026 年各引擎陆续落地)就是为此设计的:物理上是一个含两个 binary 字段的 struct(metadata 存字段名字典和类型信息,value 存紧凑二进制值),引擎可以直接定位到某个字段而不解析整块。配合 shredding,还能把高频访问的字段抽成独立的强类型列,让它重新享受编码、压缩和统计——官方博客点名的两个典型场景,一个是事件分析,另一个就是IoT 传感器数据。
Geometry:楼宇空间数据的原生类型
2026 年 2 月 Parquet 加入了原生地理空间类型(GEOMETRY,以 WKB 格式存储)。对 Pivot 的意义在于:如果将来 BIM / 空间数据(楼层轮廓、设备坐标)要落到分析层,不必再自己编码坐标——这条线和 BIM 课程会合。目前了解它存在即可。
你现在可以做的决策清单
| 决策点 | 本课程的默认建议 | 依据 |
|---|---|---|
| 冷热分界 | 固定时间规则(如 90 天) | 可预测、可自动化 |
| 分区键 | day(日增量 < 100 MB 则用 month);多租户加 tenant_id | 第 5 课:每分区 100 MB–1 GB |
| 排序键 | device_id, ts(时间维度已由分区承担) | 第 4 课:跳过率 1/44 vs 44/44 |
| 压缩 | zstd 默认级别 | 第 3 课:比 snappy 小 34%,写入只慢 12% |
| 编码 | ts 用 DELTA_BINARY_PACKED,字典收窄到低基数列 | 第 3 课:该列 7.85 → 0.03 MB |
row_group_size | 200,000 行 | 第 5 课:兼顾跳过粒度与字典有效性 |
| Page Index | 开 | 第 4 课:代价极小的第三层跳过 |
| Bloom filter | 不开(除非有高基数乱序列的点查) | 第 4 课:实测在排序列上零收益 |
| 表格式 | 先不上,用裸目录 + 幂等作业 | 本课三信号判据 |
| 退路 | 删源延后 7 天;保留导出脚本与校验记录 | 归档不可逆,必须留窗口 |
用 Docker 起一个一次性 PG,灌数据,导出,校验,删容器。全程不碰你的真实库:
$ docker run -d --name pq-lab -e POSTGRES_PASSWORD=lab -e POSTGRES_DB=pivot \ -p 55432:5432 postgres:16-alpine $ docker exec -i pq-lab psql -U postgres -d pivot <<'SQL' CREATE TABLE iot_reading ( id bigserial PRIMARY KEY, ts timestamptz NOT NULL, tenant_id text, building_code text, device_id text, metric text, value double precision, quality int); INSERT INTO iot_reading (ts, tenant_id, building_code, device_id, metric, value, quality) SELECT '2026-01-01 00:00:00+00'::timestamptz + (g % 21600) * interval '1 minute', 'T-01', 'BLD-001', 'AHU-01-' || lpad(((g / 21600) + 1)::text, 3, '0'), 'power_kw', 50 + random() * 10, 0 FROM generate_series(0, 431999) g; SELECT count(*), pg_size_pretty(pg_total_relation_size('iot_reading')) FROM iot_reading; SQL
然后在 Python 里跑本课开头那条 ATTACH + COPY(把 s3://… 换成本地目录 data/archive),最后校验:
a = con.sql("SELECT count(*), sum(value) FROM pg.public.iot_reading").fetchall() b = con.sql("SELECT count(*), sum(value) FROM read_parquet('data/archive/**/*.parquet')").fetchall() print(a, b, "一致" if a == b else "不一致")
预期结果:PG 侧 51 MB / 432,000 行,导出后 15 个文件约 3.16 MB,两侧 count 与 sum 完全相等。做完记得 docker rm -f pq-lab。
Pivot 的冷数据分层 = PG 热层(事务、点查、更新)+ Parquet 冷层(只读、按天分区、内部按设备排序)+ 固定的时间分界规则。导出可以是一条 DuckDB COPY(实测 51 MB → 3.16 MB,校验和一致),但要包上幂等(临时路径提交)、校验(行数 + 求和)、删源延后这三层保护。裸目录够用到出现"并发写 / 行级删改 / 复杂 schema 演进"中的两个信号为止——在那之前别为了名字上湖仓。Variant 类型是 IoT 异构 payload 的正解方向,但生态尚在铺开,现在了解、暂不押注。
团队每天一个作业往归档目录写一次,只会新增可空列,从不改历史。有人提议上 Iceberg。最合理的回应是什么?
课程主线到这里结束,但对话不结束。你提到手上有具体任务要交付——把它的细节告诉我(哪张表、日增量多少、下游谁在查),我会用这六课的框架帮你把方案写成可执行的脚本和一份 ADR。另外,速查表是为你日常写代码时准备的,术语表用来和同事对齐说法。