
在实际项目中测试数据不足往往是压测、联调、回归和报表口径验证最容易被低估的风险。Terafab 是本文为这套数据生成流程起的示例项目代号目标是用配置驱动的方式搭建一个能够稳定产出大规模测试数据的生成工厂。它不是某个开源产品的测评也不是商业平台的介绍而是一套可由你自己实现的工程结构列生成器负责把表字段变成可配置规则任务分片负责把百万级行数切碎多进程负责并行消费落盘和 manifest 负责让产物可校验。本文会从零实现一个在本机可运行的最小版本并说明进入生产环境后还需要补齐哪些能力。1. 为什么要把测试数据生成做成“工厂”1.1 一个脚本和一套工厂的差距很多开发者的第一版测试数据脚本长这样用 for 循环生成一万行再拼成 SQL 或 CSV。这种写法在数据量小的时候没有问题但它很快就暴露几个限制改一个字段规则就要改代码。十万行时还能接受百万行时单进程明显变慢。每次执行生成的随机结果不一致问题复现困难。没有产物校验文件是否完整、记录数是否对得上全靠人工看。Terafab 采用“工厂”式设计核心思路是把“生成规则”与“执行流程”分离。你只需要在 YAML 里描述表结构和字段规则执行器负责把规则翻译成每一行的实际值。这样表结构、字段类型、数据分布、输出格式、压缩方式和并发数都变成了配置而不是散落在代码里的循环逻辑。1.2 Terafab 的四条技术主线从工程结构看这套数据工厂有四个关键点第一是规则可声明。每一列由 name、type、params 三部分定义后续把规则扩展为代码时不需要改动整体框架。第二是任务可分片。总行数按 shard_size 切成多个区间每个区间由一个 worker 独立生成写独立的文件。这样单机可以多进程并行将来上集群时也可以把每个 shard 分发给不同节点。第三是生成可复现。所有随机源都基于 seed 和 shard_id 派生同一个 YAML、同一个 seed、相同的代码版本生成结果应该一致。这对回归验证尤其重要。第四是产物可校验。除了数据文件还会生成 manifest.json记录每个 shard 的文件路径、记录数、文件 SHA-256 和生成配置摘要。后续校验时直接比对文件哈希和行数就能发现文件被串改、半途中断或者记录数不符合预期的问题。1.3 适用场景和不适用场景Terafab 这类数据工厂适合以下场景压测平台需要生成千万级用户、订单、流水数据。数据仓库联调时需要大量历史数据灌入。回归测试需要保证每次测试数据一致。多环境部署时需要快速造出结构一致的模拟数据。它不适合直接用于生产环境真实数据的迁移、同步或脱敏。真实数据脱敏涉及更复杂的血缘关系、隐私合规和字段映射策略不是一个按配置生成随机值的工具能替代的。Terafab 的定位是“制造数据”而不是“搬运和转换数据”。2. 核心设计列生成器、任务分片和产物校验2.1 列生成器把字段规则变成代码插件列生成器是 Terafab 最基础的抽象。每个字段需要什么值由type决定生成时需要的边界和候选值由params提供。例如randint生成整数参数为min和maxdate生成日期字符串参数为start和end。下面的表格给出最小版本中需要内置的生成器类型类型用途主要参数输出示例sequence自增整数start1001, 1002constant固定值valuetruerandint整数范围min,max18, 45random_float浮点数min,max,precision67.35enum从候选列表取值valuesshanghainame中文姓名无张伟phone手机号无13812345678date日期start,end2023-08-15uuid唯一标识无1f3a...把生成器做成一等公民后外层累进循环只需要遍历列逐个调用对应的生成逻辑。将来要增加身份证号、经纬度、JSON 嵌套字段等类型也只需要新增一个elif分支或注册函数不需要改动并发和落盘逻辑。2.2 任务分片让生成可以水平扩展一个简单的数据脚本通常只有一个循环从第 0 行跑到第 total_rows 行。Terafab 把这一过程拆成多个 shard例如 total_rows 为 100 万shard_size 为 10 万就得到 10 个任务shard 0: 第 0 至 99999 行 shard 1: 第 100000 至 199999 行 ...每个 shard 之间没有依赖因此可以交给进程池并发执行。分片带来的额外收益体现在两点单个文件不会过大默认按行数切分后续读取和排查都方便。如果某个 shard 失败可以只重新生成失败区间不需要重跑整个任务。这里的区间使用半开区间[start, end)也就是包含 start不包含 end。这个约定会影响最终记录数后面排错部分会专门说明。2.3 产物与 manifest数据和元数据一起产出生成完成后output_dir 下会有一组数据文件例如user_00000000.csv.gz user_00000001.csv.gz ... manifest.jsonmanifest.json 是这套设计的“收据”至少包含{ table: user, generated_at: 2025-01-01T12:30:00, config: { table: user, total_rows: 1000000, shard_size: 100000, seed: 42 }, shards: [ { shard_id: 0, file: data/user_00000000.csv.gz, records: 100000, sha256: a4f5... } ], total_records: 1000000 }这样即使数据文件被误删或改动也可以基于 manifest 做校验。生产环境甚至可以把 manifest 写入 Hive 或数据库的分区元数据作为数据表的一部分。2.4 为什么用 YAML 描述结构YAML 相比代码更接近表结构本身阅读成本低适合给非研发同学评审。它天然支持列表和嵌套能比较直观地表达一个表有哪些列、每列是什么类型。不过 YAML 也有一个需要警惕的地方它的类型推断规则比较宽松。start: 2023-01-01在 PyYAML 中可能被解析成date对象而不是字符串。所以配置加载后最好人工做好默认值和类型处理。下面的实现中date 生成器会通过strptime显式解析字符串避免不同 YAML 解析结果带来的不确定性。3. 环境准备与目录结构3.1 运行环境依赖这个最小版本使用 Python 3.10 或更高版本第三方依赖只有一个 PyYAML。生成、并发、压缩、哈希和校验全部使用标准库实现。安装依赖前先确认 Python 版本python --version然后创建虚拟环境并安装依赖python -m venv .venv source .venv/bin/activate pip install pyyaml在 Windows 上激活虚拟环境的命令是.venv\Scripts\activate如果想使用 JSON 格式配置可以把 PyYAML 去掉但 YAML 的可读性更适合作为示例。3.2 建议目录结构最小项目可以放在一个目录下terafab/ config.yaml terafab.py data/config.yaml描述表结构和执行参数terafab.py是核心脚本data/用于存放产物。这样设计的好处是脚本本身不依赖外部服务跑一次就能看到完整链路。3.3 配置字段速查字段是否必填含义建议值table是表名也是输出文件前缀usercolumns是列生成规则列表至少一列total_rows是总生成行数根据场景调整shard_size否每个 shard 行数默认 10000010000 至 100000concurrency否并发进程数默认 CPU 核数小机器从 2 开始output_dir否输出目录默认 datadataformat否csv 或 jsonl默认 csvcsvcompression否none 或 gzip默认 nonegzipseed否随机种子默认 42固定整数需要提醒的是concurrency并不是越大越好。在小文件场景中瓶颈往往是磁盘写入而不是 CPU。并发太高会导致大量 worker 同时写文件反而增加 IO 抖动。4. 核心代码实现4.1 配置加载与结构校验配置加载不只是读 YAML还要做基本校验避免错误配置走到执行阶段才暴露出问题。下面这段代码定义了必填项和默认值import json import os import yaml from pathlib import Path def load_config(path: Path) - dict: cfg yaml.safe_load(path.read_text(encodingutf-8)) required {table, columns, total_rows} missing required - set(cfg.keys()) if missing: raise ValueError(fconfig missing: {missing}) total_rows int(cfg[total_rows]) if total_rows 0: raise ValueError(total_rows must be positive) if not isinstance(cfg[columns], list) or len(cfg[columns]) 0: raise ValueError(columns must be a non-empty list) cfg.setdefault(output_dir, data) cfg.setdefault(shard_size, 100000) cfg.setdefault(concurrency, os.cpu_count() or 4) cfg.setdefault(format, csv) cfg.setdefault(compression, none) cfg.setdefault(seed, 42) return cfg这里的关键是“失败要早”如果配置缺少 table 或 total_rows主程序应该直接退出而不是生成一堆内容后再发现文件名和行数不对。setdefault用来补齐默认值让后续执行代码不需要反复判断字段是否存在。4.2 内置生成器与单行生成在最小实现中列生成逻辑用一个函数完成即可。它接收配置、当前行号和随机数对象返回一行数据import random from datetime import datetime, timedelta FIRST_NAMES [王, 李, 张, 刘, 陈, 杨, 赵, 黄] LAST_NAMES [伟, 芳, 娜, 敏, 磊, 洋, 勇, 军, 杰, 婷] def generate_row(cfg: dict, row_index: int, rng: random.Random) - list: row [] for col in cfg[columns]: col_type col[type] params col.get(params, {}) or {} if col_type sequence: value params.get(start, 0) row_index elif col_type constant: value params.get(value) elif col_type randint: value rng.randint(int(params.get(min, 0)), int(params.get(max, 100))) elif col_type random_float: value round( rng.uniform(float(params.get(min, 0.0)), float(params.get(max, 1.0))), int(params.get(precision, 2)), ) elif col_type enum: values params.get(values, []) value values[rng.randrange(len(values))] elif col_type name: value rng.choice(FIRST_NAMES) rng.choice(LAST_NAMES) elif col_type phone: value 1 .join(str(rng.randint(0, 9)) for _ in range(10)) elif col_type date: start_dt datetime.strptime(params.get(start, 2023-01-01), %Y-%m-%d) end_dt datetime.strptime(params.get(end, 2024-12-31), %Y-%m-%d) days max((end_dt - start_dt).days, 0) value (start_dt timedelta(daysrng.randint(0, days))).date().isoformat() elif col_type uuid: value -.join([ .join(rng.choice(0123456789abcdef) for _ in range(8)), .join(rng.choice(0123456789abcdef) for _ in range(4)), 4 .join(rng.choice(0123456789abcdef) for _ in range(3)), rng.choice(89ab) .join(rng.choice(0123456789abcdef) for _ in range(3)), .join(rng.choice(0123456789abcdef) for _ in range(12)), ]) else: raise ValueError(funsupported type: {col_type}) row.append(value) return row这里使用外部传入的rng而不是在函数内部调用random.randint。原因是 Python 全局随机函数的内存随机数状态会让结果不可控。外部传入random.Random(seed)后只要 seed 相同生成序列就完全一致。4.3 分片生成与多进程执行分片逻辑有两种写法普通循环和ProcessPoolExecutor。多进程实现的核心是把每个 shard 的参数封装成 task再由进程池调度from concurrent.futures import ProcessPoolExecutor from pathlib import Path def shard_ranges(total_rows: int, shard_size: int): start 0 while start total_rows: end min(total_rows, start shard_size) yield start, end start end def build_tasks(cfg: dict): tasks [] for shard_id, (start, end) in enumerate(shard_ranges(cfg[total_rows], cfg[shard_size])): tasks.append({ shard_id: shard_id, start: start, end: end, config: cfg, }) return tasks真正的写文件逻辑放在generate_shard函数中。它会被多个进程调用因此必须不依赖全局可变状态import csv import gzip import hashlib def open_text(path, compression): if compression gzip: return gzip.open(path, wt, encodingutf-8, newline) return open(path, w, encodingutf-8, newline) def sha256_file(path): h hashlib.sha256() with open(path, rb) as f: for chunk in iter(lambda: f.read(1024 * 1024), b): h.update(chunk) return h.hexdigest() def generate_shard(task: dict) - dict: cfg task[config] shard_id task[shard_id] start task[start] end task[end] rng random.Random(f{cfg[seed]}-{shard_id}) table cfg[table] output_dir Path(cfg[output_dir]) file_ext .csv if cfg[format] csv else .jsonl if cfg[compression] gzip: file_ext .gz file_path output_dir / f{table}_{shard_id:08d}{file_ext} header [col[name] for col in cfg[columns]] if cfg[format] csv: with open_text(file_path, cfg[compression]) as f: writer csv.writer(f, lineterminator\n) writer.writerow(header) for idx in range(start, end): writer.writerow(generate_row(cfg, idx, rng)) else: with open_text(file_path, cfg[compression]) as f: for idx in range(start, end): line dict(zip(header, generate_row(cfg, idx, rng))) f.write(json.dumps(line, ensure_asciiFalse) \n) return { shard_id: shard_id, file: str(file_path), records: end - start, sha256: sha256_file(file_path), } def run_tasks(cfg: dict): output_dir Path(cfg[output_dir]) output_dir.mkdir(parentsTrue, exist_okTrue) with ProcessPoolExecutor(max_workerscfg[concurrency]) as pool: results list(pool.map(generate_shard, build_tasks(cfg))) manifest { table: cfg[table], generated_at: datetime.now().isoformat(timespecseconds), config: cfg, shards: results, total_records: sum(item[records] for item in results), } manifest_path output_dir / manifest.json manifest_path.write_text(json.dumps(manifest, ensure_asciiFalse, indent2), encodingutf-8) return manifest_path注意generate_shard中所有内容都基于局部变量和传入的 task没有读写模块级全局数组。这是多进程代码是否稳定的分水岭。Windows 上使用进程池时入口处必须加if __name__ __main__否则新进程会重复执行模块顶层代码。4.4 落盘、压缩与 manifest 的关系整个执行流程中落盘和 manifest 是同一阶段的两个动作。先写数据文件再计算 SHA-256最后把所有 shard 信息汇总成 manifest.json。顺序不能反过来因为 manifest 里的哈希必须在文件写完并关闭后计算。计算哈希时按 1MB 块读取文件而不是一次性read()这样可以避免大文件把内存占满。manifest 中保存的文件路径建议使用相对路径方便把整个 output_dir 迁移到其他机器。如果使用绝对路径换服务器后校验脚本会因为路径不一致而失败。5. 运行与验证5.1 准备一个演示配置假设要生成一张千万级用户表字段包括自增 id、姓名、年龄、城市、注册时间和分数。配置如下table: user columns: - name: id type: sequence params: start: 1 - name: user_name type: name - name: age type: randint params: min: 18 max: 60 - name: city type: enum params: values: [beijing, shanghai, shenzhen, chengdu] - name: register_date type: date params: start: 2023-01-01 end: 2024-12-31 - name: score type: random_float params: min: 0.0 max: 100.0 precision: 2 total_rows: 1000000 shard_size: 100000 concurrency: 4 output_dir: data format: csv compression: gzip seed: 42这个配置的意思是生成 100 万行每 10 万行一个 shard4 个进程并发。唯一能保证一致性的随机来源是 seed所以 seed 需要固定。5.2 启动生成命令行运行方式可以是python terafab.py run --config config.yaml正常结束后data 目录下会有 10 个 gzip 压缩的 CSV 文件以及 manifest.json。可以在终端确认文件ls -lh data预期看到类似输出user_00000000.csv.gz user_00000001.csv.gz ... user_00000009.csv.gz manifest.json5.3 校验生成产物如果实现了 validate 子命令可以这样校验python terafab.py validate --manifest data/manifest.json校验逻辑会依次检查每个 shard 文件是否缺失、文件哈希是否与 manifest 一致、实际行数是否等于 records。只要其中一项不匹配就以非零状态退出。这个命令应该被放进 CI 或数据接入任务中确保生成环节没有静默失败。5.4 单机小压测参数单机压测时可以先从下面这组参数看效果参数最小值建议值说明total_rows10001000000先小规模验证 schemashard_size1000100000太小会导致文件碎片过多concurrency1CPU 核数磁盘 IO 瓶颈时降低compressionnonegzip文件越小后续传输越快先跑小规模配置再放大行数。不要一上来就压千万行否则很容易把磁盘写满而配置里的字段类型错误还没发现。6. 常见问题和排查路径6.1 生成的记录数不等于 total_rows现象是 manifest 中 total_records 比配置少或者多出若干行。常见原因是区间边界处理不一致。如果使用range(start, end)Python 的 range 不包含 end所以总行数一定是end - start。若在使用过程中把 end 改成了end 1就会多生成一行若把 total_rows 除以 shard_size 后只取整数商就会丢掉最后一个不足 shard_size 的余数段。检查方式是看任务构建代码list(shard_ranges(1000000, 100000)) # 应该有 10 个区间最后一个区间是 [900000, 1000000)预防措施是在 manifest 中输出 total_records并把 validate 命令纳入常规流程。6.2 Windows 下多进程报错或重复执行现象是执行 run 命令时程序抛错错误信息指向ProcessPoolExecutor或者模块顶层代码被执行了多次。原因是在 Windows 上进程池使用 spawn 方式创建子进程子进程会重新导入主模块。如果没有if __name__ __main__:保护模块顶层的命令会被重复执行。处理方式很明确把命令行入口放到main()再在模块末尾加if __name__ __main__: main()所有全局初始化也应放在if __name__ __main__块内而不是模块顶层。6.3 随机数不可复现前后两次数据不一致现象是同一份 YAML、同一个 seed两次生成出来的内容不同。原因是代码中混用了全局random模块和自定义random.Random。全局随机函数状态由解释器维护运行环境和调用顺序都会影响结果。只用全局randomseed 只能保证上下文一致时稳定一旦前面多调用了一次random后续结果就全变了。推荐做法是为每个 shard 单独创建rng random.Random(f{seed}-{shard_id})并使用这个rng完成该 shard 内所有随机调用。这样每个 shard 的序列只依赖 seed 和 shard_id不受其他 shard 影响。6.4 CSV 文件在 Excel 中打开中文乱码现象是生成的数据用文本编辑器正常用 Excel 打开后中文显示乱码。原因是 CSV 默认使用 UTF-8 编码而老版本 Excel 默认按 GBK 或本地区域编码解析。如果只想保证 Excel 兼容可以把 CSV 写入编码改为utf-8-sigopen(path, w, encodingutf-8-sig, newline)但要注意如果后续数据要进入 Linux 下的数据处理流程utf-8-sig会在文件开头多一个 BOM 标记部分解析器需要额外处理。建议根据数据消费方选择编码不要一刀切。