""" Parquet 课程 · 样本数据生成器(全课程共用) 生成一份 Pivot 风格的楼宇 IoT 时序样本:200 台设备 × 15 天 × 每分钟 1 点 = 4,320,000 行, 并落成 CSV / JSONL / Parquet 四种文件,供第 1–5 课的动手环节使用。 pip install pyarrow duckdb # 只需要这两个 python gen_sample.py # 约 1 分钟,产出 ./data/ 课程正文里的每一个体积数字,都出自这个脚本 + 本机实测(pyarrow 25.0.0 / DuckDB 1.5.5)。 """ import os, json, random import pyarrow as pa import pyarrow.parquet as pq import pyarrow.csv as pacsv random.seed(42) # 固定随机种子,你跑出来的数字应与课程一致 OUT = os.path.join(os.path.dirname(os.path.abspath(__file__)), "data") os.makedirs(OUT, exist_ok=True) N_DEV, DAYS, STEP_S = 200, 15, 60 ROWS_PER_DEV = DAYS * 24 * 3600 // STEP_S METRICS = ["power_kw", "temp_c", "humidity", "co2_ppm", "valve_pct", "flow_m3h"] BASE_TS = 1767225600_000 # 2026-01-01T00:00:00Z,毫秒 ts, dev, met, val, qua, bld, ten = [], [], [], [], [], [], [] for d in range(N_DEV): device = f"AHU-{d//20+1:02d}-{d%20+1:03d}" metric = METRICS[d % len(METRICS)] building = f"BLD-{d//50+1:03d}" tenant = f"T-{d//100+1:02d}" v = 50.0 for i in range(ROWS_PER_DEV): ts.append(BASE_TS + i * STEP_S * 1000) dev.append(device); met.append(metric) bld.append(building); ten.append(tenant) v += random.gauss(0, 0.35) # 缓慢游走,模拟真实传感器 val.append(round(v, 3)) qua.append(0 if random.random() > 0.01 else 1) # 1% 坏点 tbl = pa.table({ "ts": pa.array(ts, pa.timestamp("ms", tz="UTC")), "tenant_id": pa.array(ten, pa.string()), "building_code": pa.array(bld, pa.string()), "device_id": pa.array(dev, pa.string()), "metric": pa.array(met, pa.string()), "value": pa.array(val, pa.float64()), "quality": pa.array(qua, pa.int32()), }) print(f"生成 {tbl.num_rows:,} 行 / {tbl.num_columns} 列") # 打散:模拟"按 MQTT 到达顺序"写入,而不是理好序再写 shuffled = tbl.take(pa.array(random.sample(range(tbl.num_rows), tbl.num_rows))) pacsv.write_csv(shuffled, f"{OUT}/raw.csv") cols = shuffled.to_pydict() with open(f"{OUT}/raw.jsonl", "w") as f: for i in range(tbl.num_rows): f.write(json.dumps({k: (v[i].isoformat() if k == "ts" else v[i]) for k, v in cols.items()}, separators=(",", ":")) + "\n") pq.write_table(shuffled, f"{OUT}/iot_unsorted.parquet", compression="zstd") pq.write_table(shuffled.sort_by([("device_id", "ascending"), ("ts", "ascending")]), f"{OUT}/iot_sorted.parquet", compression="zstd") # 第 3 课要用:关掉全部编码与压缩的"裸"版本 pq.write_table(shuffled, f"{OUT}/iot_naked.parquet", compression="none", use_dictionary=False) for f in sorted(os.listdir(OUT)): p = f"{OUT}/{f}" if os.path.isfile(p): print(f"{f:26s} {os.path.getsize(p)/1024/1024:9.2f} MB")