databricks-synthetic-data-gen
Compare original and translation side by side
🇺🇸
Original
English🇨🇳
Translation
ChineseCatalog and schema are always user-supplied — never default to any value. If the user hasn't provided them, ask. For any UC write, always create the schema if it doesn't exist before writing data.
Catalog和schema始终由用户提供——绝不要使用默认值。如果用户未提供,需主动询问。对于任何UC写入操作,在写入数据前请务必先创建不存在的schema。
Databricks Synthetic Data Generation
Databricks合成数据生成
Generate realistic, story-driven synthetic data for Databricks using Spark + Faker + Pandas UDFs (strongly recommended).
使用Spark + Faker + Pandas UDFs(强烈推荐)为Databricks生成逼真的、带有业务场景的合成数据。
Data Must Tell a Business Story
数据需承载业务场景
Synthetic data should demonstrate how Databricks helps solve real business problems.
The pattern: Something goes wrong → business impact ($) → analyze root cause → identify affected customers → fix and prevent.
Key principles:
- Problem → Impact → Analysis → Solution — Include an incident, anomaly, or issue that causes measurable business impact. The data lets you find the root cause and act on it.
- Industry-relevant but simple — Use domain terms (e.g., "SLA breach", "churn", "stockout") but keep the schema easy to understand. A few tables, clear relationships.
- Business metrics with $ impact — Revenue, MRR, cost, conversion rate. Every story needs a dollar sign to show why it matters.
- Tables explain each other — Ticket spike? Incident table shows the outage. Revenue drop? Churn table shows who left and why. All data connects.
- Actionable insights — Data should answer: What happened? Who's affected? How much did it cost? How do we prevent it?
Why no flat distributions: Uniform data has no story — no spikes, no anomalies, no cohort, no 20/80, no skew, nothing to investigate. It can't show Databricks' value for root cause analysis.
合成数据应展示Databricks如何帮助解决实际业务问题。
模式: 出现问题 → 业务影响(金额) → 分析根本原因 → 识别受影响客户 → 修复并预防。
核心原则:
- 问题→影响→分析→解决方案——包含会造成可衡量业务影响的事件、异常或问题。数据需能帮助找到根本原因并采取行动。
- 贴合行业但简洁易懂——使用领域术语(如“SLA违约”“客户流失”“库存短缺”),但保持schema易于理解。仅需少量表,关系清晰。
- 带金额影响的业务指标——收入、MRR、成本、转化率。每个场景都需要体现金额,以说明其重要性。
- 表之间相互关联——工单激增?事件表显示系统中断。收入下降?流失表显示客户离开的原因。所有数据相互关联。
- 可落地的洞察——数据应能回答:发生了什么?谁受影响?造成了多少损失?如何预防?
为何不使用均匀分布: 均匀数据没有场景——没有峰值、没有异常、没有群体差异、没有20/80法则、没有偏斜,没有可调查的内容。无法展示Databricks在根本原因分析中的价值。
References
参考资料
| When | Guide |
|---|---|
| User mentions ML model training or complex time patterns | references/1-data-patterns.md — ML-ready data, time multipliers, row coherence |
| Errors during generation | references/2-troubleshooting.md — Fixing common issues |
| 场景 | 指南 |
|---|---|
| 用户提及ML模型训练或复杂时间模式 | references/1-data-patterns.md — 适用于ML的数据、时间乘数、行一致性 |
| 生成过程中出现错误 | references/2-troubleshooting.md — 修复常见问题 |
Critical Rules
关键规则
- Data tells a story — Something goes wrong, impacts $, can be analyzed and fixed. Show Databricks value.
- All data serves the story — Every table and column must be coherent and usable in dashboards or ML models. No orphan data, no random noise — if it doesn't help explain or plot a futur dashboard or predict, don't generate it.
- Industry terms, simple schema — Use domain-specific vocabulary but keep it easy to understand (few tables, clear relationships)
- Never uniform distributions — Skewed categories, log-normal amounts, 80/20 patterns. Flat = no story = useless
- Enough data for trends — ~100K+ rows for main tables so patterns survive aggregation
- Ask for catalog/schema — Never default, always confirm before generating
- Present plan for approval — Show tables, distributions, assumptions before writing code
- Master tables first — Generate parent tables, write to Delta, then create children with valid FKs
- Use Spark + Faker + Pandas UDFs — Scalable, parallel. Polars only if user explicitly wants local + <30K rows
- Use Databricks Connect Serverless by default to generate data — Update databricks-connect on python 3.12 if required (avoid using execute_code unless instructed to not use Databricks Connect)
- No or
.cache()— Not supported on serverless. Write to Delta, read back for joins.persist() - No Python loops or — Use Spark parallelism. No driver-side iteration, avoid Pandas↔Spark conversions
.collect()
- 数据承载场景——出现问题,影响金额,可分析并修复。展示Databricks的价值。
- 所有数据服务于场景——每张表和每列都必须连贯,可用于仪表板或ML模型。无孤立数据,无随机噪声——如果无法帮助解释、绘制未来仪表板或进行预测,则不要生成。
- 行业术语,简洁schema——使用领域特定词汇,但保持易懂(少量表,关系清晰)
- 绝不使用均匀分布——偏斜分类、对数正态金额、80/20模式。均匀分布=无场景=无用
- 足够数据以呈现趋势——主表需约10万+行,确保聚合后模式仍存在
- 询问Catalog/Schema——绝不默认,生成前务必确认
- 提交计划供批准——编写代码前展示表、分布、假设
- 先生成主表——生成父表,写入Delta,再创建带有有效外键的子表
- 使用Spark + Faker + Pandas UDFs——可扩展、并行化。仅当用户明确要求本地生成且行数<30K时使用Polars
- 默认使用Databricks Connect Serverless生成数据——必要时在Python 3.12上更新databricks-connect(除非指示不使用Databricks Connect,否则避免使用execute_code)
- 禁止使用或
.cache()——无服务器环境不支持。写入Delta后再读取进行关联.persist() - 禁止Python循环或——使用Spark并行化。避免驱动端迭代,减少Pandas与Spark之间的转换
.collect()
Generation Planning Workflow
生成规划流程
Before generating any code, you MUST present a plan for user approval.
在生成任何代码之前,必须提交计划供用户批准。
⚠️ MUST DO: Confirm Catalog Before Proceeding
⚠️ 必须执行:先确认Catalog
You MUST explicitly ask the user which catalog to use. Do not assume or proceed without confirmation.
Example prompt to user:
"Which Unity Catalog should I use for this data?"
When presenting your plan, always show the selected catalog prominently:
📍 Output Location: catalog_name.schema_name
Volume: /Volumes/catalog_name/schema_name/raw_data/This makes it easy for the user to spot and correct if needed.
必须明确询问用户使用哪个Catalog。 不得假设或未确认就继续。
示例用户提示:
“我应该为这些数据使用哪个Unity Catalog?”
提交计划时,务必突出显示所选Catalog:
📍 输出位置: catalog_name.schema_name
卷: /Volumes/catalog_name/schema_name/raw_data/这样用户可以轻松发现并纠正错误。
Step 1: Gather Requirements
步骤1:收集需求
Ask the user about:
- Catalog/Schema — Which catalog to use?
- Domain — E-commerce, support tickets, IoT, financial? (Use industry terms)
If user doesn't specify a story: Propose one. Don't generate bland data — suggest an incident, anomaly, or trend that shows Databricks value (e.g., "I'll include a system outage that causes ticket spike and churn — this lets you demo root cause analysis").
询问用户以下信息:
- Catalog/Schema——使用哪个Catalog?
- 领域——电商、支持工单、IoT、金融?(使用行业术语)
如果用户未指定场景: 主动提议。不要生成平淡的数据——建议一个能展示Databricks价值的事件、异常或趋势(例如:“我将包含一个导致工单激增和客户流失的系统中断场景——这可以演示根本原因分析”)。
Step 2: Present Plan with Story
步骤2:提交带场景的计划
Show a clear specification with the business story and your assumptions surfaced:
📍 Output Location: {user_catalog}.support_demo
Volume: /Volumes/{user_catalog}/support_demo/raw_data/
📖 Story: A payment system outage causes support ticket spike. Resolution times
degrade, enterprise customers churn, revenue drops $2.3M. With Databricks we
identify the root cause, affected customers, and prevent future impact.| Table | Description | Rows | Key Assumptions |
|---|---|---|---|
| customers | Customer profiles with tier, MRR | 10,000 | Enterprise 10% but 60% of revenue |
| tickets | Support tickets with priority, resolution_time | 80,000 | Spike during outage, SLA breaches |
| incidents | System events (outages, deployments) | 50 | Payment outage mid-month |
| churn_events | Customer cancellations with reason | 500 | Spike after poor support experience |
Business metrics:
- — Revenue at risk ($)
customers.mrr - — SLA performance
tickets.resolution_hours - — Churn impact ($)
churn_events.lost_mrr
The story this data tells:
- Incident table shows payment outage on March 15
- Tickets spike 5x during outage, resolution time degrades from 4h → 18h
- Enterprise customers with SLA breaches churn 3 weeks later
- Total impact: $2.3M lost MRR, traceable to one incident
- Databricks value: Root cause analysis, identify at-risk customers, build alerting
Ask user: "Does this story work? Any adjustments?"
展示清晰的规范,明确业务场景和你的假设:
📍 输出位置: {user_catalog}.support_demo
卷: /Volumes/{user_catalog}/support_demo/raw_data/
📖 场景:支付系统中断导致支持工单激增。解决时间变长,企业客户流失,收入减少230万美元。借助Databricks,我们可以识别根本原因、受影响客户,并防止未来的影响。| 表 | 描述 | 行数 | 核心假设 |
|---|---|---|---|
| customers | 包含客户等级、MRR的客户档案 | 10,000 | 企业客户占10%但贡献60%的收入 |
| tickets | 包含优先级、解决时间的支持工单 | 80,000 | 中断期间工单激增,出现SLA违约 |
| incidents | 系统事件(中断、部署) | 50 | 月中发生支付系统中断 |
| churn_events | 包含原因的客户取消记录 | 500 | 糟糕的支持体验后客户流失激增 |
业务指标:
- — 面临风险的收入(美元)
customers.mrr - — SLA表现
tickets.resolution_hours - — 客户流失影响(美元)
churn_events.lost_mrr
此数据讲述的场景:
- 事件表显示3月15日发生支付系统中断
- 中断期间工单激增5倍,解决时间从4小时延长至18小时
- 遭遇SLA违约的企业客户在3周后流失
- 总影响:损失230万美元MRR,可追溯至单个事件
- Databricks价值: 根本原因分析、识别高风险客户、构建告警系统
询问用户:“这个场景是否合适?需要调整吗?”
Step 3: Ask About Data Features
步骤3:询问数据特性
- Skew (non-uniform distributions) - Enabled by default
- Joins (referential integrity) - Enabled by default
- Bad data injection (for data quality testing)
- Multi-language text
- Incremental mode (append instead of overwrite)
- 偏斜分布(非均匀)- 默认启用
- 关联(引用完整性)- 默认启用
- 注入坏数据(用于数据质量测试)
- 多语言文本
- 增量模式(追加而非覆盖)
Pre-Generation Checklist
生成前检查清单
- Catalog confirmed - User explicitly approved which catalog to use
- Output location shown prominently in plan (easy to spot/change)
- Table specification shown and approved
- Assumptions about distributions confirmed
- User confirmed compute preference (Databricks Connect on serverless recommended)
- Data features selected
Do NOT proceed to code generation until user approves the plan, including the catalog.
- 已确认Catalog - 用户明确批准使用的Catalog
- 计划中突出显示输出位置(易于发现/修改)
- 表规范已展示并获批
- 分布假设已确认
- 用户确认计算偏好(推荐使用无服务器Databricks Connect)
- 已选择数据特性
在用户批准计划(包括Catalog)之前,不得进行代码生成。
Post-Generation Validation
生成后验证
Use to validate generated data (row counts, distributions, referential integrity). Query parquet files directly:
databricks experimental aitools tools querybash
databricks experimental aitools tools query --warehouse $WAREHOUSE_ID "
SELECT COUNT(*) FROM parquet.\`/Volumes/CATALOG/SCHEMA/raw_data/customers\`
"See references/2-troubleshooting.md for full validation examples.
使用验证生成的数据(行数、分布、引用完整性)。直接查询parquet文件:
databricks experimental aitools tools querybash
databricks experimental aitools tools query --warehouse $WAREHOUSE_ID "
SELECT COUNT(*) FROM parquet.\`/Volumes/CATALOG/SCHEMA/raw_data/customers\`
"完整验证示例请参见references/2-troubleshooting.md。
Use Databricks Connect Spark + Faker Pattern
使用Databricks Connect Spark + Faker模式
python
from databricks.connect import DatabricksSession
from pyspark.sql import functions as F
from pyspark.sql.types import StringType
import pandas as pdpython
from databricks.connect import DatabricksSession
from pyspark.sql import functions as F
from pyspark.sql.types import StringType
import pandas as pdSetup serverless Spark session
Setup serverless Spark session
spark = DatabricksSession.builder.serverless(True).getOrCreate()
spark = DatabricksSession.builder.serverless(True).getOrCreate()
Pandas UDF pattern - import lib INSIDE the function (libs must be installed locally)
Pandas UDF pattern - import lib INSIDE the function (libs must be installed locally)
@F.pandas_udf(StringType())
def fake_name(ids: pd.Series) -> pd.Series:
from faker import Faker # Import inside UDF
fake = Faker()
return pd.Series([fake.name() for _ in range(len(ids))])
@F.pandas_udf(StringType())
def fake_name(ids: pd.Series) -> pd.Series:
from faker import Faker # Import inside UDF
fake = Faker()
return pd.Series([fake.name() for _ in range(len(ids))])
Generate with spark.range, apply UDFs
Generate with spark.range, apply UDFs
customers_df = spark.range(0, 10000, numPartitions=16).select(
F.concat(F.lit("CUST-"), F.lpad(F.col("id").cast("string"), 5, "0")).alias("customer_id"),
fake_name(F.col("id")).alias("name"),
)
customers_df = spark.range(0, 10000, numPartitions=16).select(
F.concat(F.lit("CUST-"), F.lpad(F.col("id").cast("string"), 5, "0")).alias("customer_id"),
fake_name(F.col("id")).alias("name"),
)
Write to Volume as Parquet (default for raw data)
Write to Volume as Parquet (default for raw data)
Path is a folder with table name: /Volumes/catalog/schema/raw_data/customers/
Path is a folder with table name: /Volumes/catalog/schema/raw_data/customers/
spark.sql(f"CREATE SCHEMA IF NOT EXISTS {CATALOG}.{SCHEMA}")
spark.sql(f"CREATE VOLUME IF NOT EXISTS {CATALOG}.{SCHEMA}.raw_data")
customers_df.write.mode("overwrite").parquet(f"/Volumes/{CATALOG}/{SCHEMA}/raw_data/customers")
**Partitions by scale:** `spark.range(N, numPartitions=P)`
- <100K rows: 8 partitions
- 100K-500K: 16 partitions
- 500K-1M: 32 partitions
- 1M+: 64+ partitions
**Output formats:**
- **Parquet to Volume** (default): `df.write.parquet("/Volumes/.../raw_data/table")` — raw data for pipelines
- **Delta Table**: `df.write.saveAsTable("catalog.schema.table")` — if user wants queryable tables
- **JSON/CSV**: small dimension tables, replicate legacy systemsspark.sql(f"CREATE SCHEMA IF NOT EXISTS {CATALOG}.{SCHEMA}")
spark.sql(f"CREATE VOLUME IF NOT EXISTS {CATALOG}.{SCHEMA}.raw_data")
customers_df.write.mode("overwrite").parquet(f"/Volumes/{CATALOG}/{SCHEMA}/raw_data/customers")
**按规模设置分区:** `spark.range(N, numPartitions=P)`
- <10万行:8个分区
- 10万-50万行:16个分区
- 50万-100万行:32个分区
- 100万+行:64+个分区
**输出格式:**
- **Parquet写入卷**(默认):`df.write.parquet("/Volumes/.../raw_data/table")` — 用于数据管道的原始数据
- **Delta Table**:`df.write.saveAsTable("catalog.schema.table")` — 如果用户需要可查询的表
- **JSON/CSV**:小型维度表,复制遗留系统Performance Rules
性能规则
Generated scripts must be highly performant. Never do these:
| Anti-Pattern | Why It's Slow | Do This Instead |
|---|---|---|
| Python loops on driver | Single-threaded, no parallelism | Use |
| Brings all data to driver memory | Keep data in Spark, use DataFrame ops |
| Pandas → Spark → Pandas | Serialization overhead, defeats distribution | Stay in Spark, use |
| Read/write temp files | Unnecessary I/O | Chain DataFrame transformations |
| Scalar UDFs | Row-by-row processing | Use |
Good pattern: → Spark transforms → for Faker → write directly
spark.range()pandas_udf生成的脚本必须具备高性能。绝对禁止以下操作:
| 反模式 | 为何缓慢 | 替代方案 |
|---|---|---|
| 驱动端Python循环 | 单线程,无并行性 | 使用 |
| 将所有数据加载到驱动端内存 | 数据保留在Spark中,使用DataFrame操作 |
| Pandas → Spark → Pandas | 序列化开销,失去分布式优势 | 保持在Spark环境中,仅在UDF中使用 |
| 读写临时文件 | 不必要的I/O | 链式DataFrame转换 |
| 标量UDF | 逐行处理 | 使用 |
推荐模式: → Spark转换 → 调用Faker → 直接写入
spark.range()pandas_udfCommon Patterns
常见模式
Weighted Categories (never uniform)
加权分类(绝不均匀)
python
F.when(F.rand() < 0.6, "Free").when(F.rand() < 0.9, "Pro").otherwise("Enterprise")python
F.when(F.rand() < 0.6, "Free").when(F.rand() < 0.9, "Pro").otherwise("Enterprise")Log-Normal Amounts (in a pandas UDF)
对数正态金额(在pandas UDF中)
Use — always positive, long tail:
np.random.lognormal(mean, sigma)- Enterprise: → ~$1800 median
lognormal(7.5, 0.8) - Pro: → ~$245 median
lognormal(5.5, 0.7) - Free: → ~$55 median
lognormal(4.0, 0.6)
使用 — 始终为正,长尾分布:
np.random.lognormal(mean, sigma)- 企业客户:→ 中位数约$1800
lognormal(7.5, 0.8) - Pro客户:→ 中位数约$245
lognormal(5.5, 0.7) - 免费客户:→ 中位数约$55
lognormal(4.0, 0.6)
Date Range (Last 6 Months)
日期范围(过去6个月)
python
END_DATE = datetime.now()
START_DATE = END_DATE - timedelta(days=180)python
END_DATE = datetime.now()
START_DATE = END_DATE - timedelta(days=180)Infrastructure (always create in script)
基础设施(始终在脚本中创建)
python
spark.sql(f"CREATE SCHEMA IF NOT EXISTS {CATALOG}.{SCHEMA}")
spark.sql(f"CREATE VOLUME IF NOT EXISTS {CATALOG}.{SCHEMA}.raw_data")python
spark.sql(f"CREATE SCHEMA IF NOT EXISTS {CATALOG}.{SCHEMA}")
spark.sql(f"CREATE VOLUME IF NOT EXISTS {CATALOG}.{SCHEMA}.raw_data")Referential Integrity (FK pattern)
引用完整性(外键模式)
Write master table to Delta first, then read back for FK joins (no on serverless):
.cache()python
undefined先将主表写入Delta,再读取回来进行外键关联(无服务器环境禁止使用):
.cache()python
undefined1. Write master table
1. 写入主表
customers_df.write.mode("overwrite").saveAsTable(f"{CATALOG}.{SCHEMA}.customers")
customers_df.write.mode("overwrite").saveAsTable(f"{CATALOG}.{SCHEMA}.customers")
2. Read back for FK lookup
2. 读取回来用于外键查找
customer_lookup = spark.table(f"{CATALOG}.{SCHEMA}.customers").select("customer_idx", "customer_id")
customer_lookup = spark.table(f"{CATALOG}.{SCHEMA}.customers").select("customer_idx", "customer_id")
3. Generate child table with valid FKs via join
3. 通过关联生成带有有效外键的子表
orders_df = spark.range(N_ORDERS).select(
(F.abs(F.hash(F.col("id"))) % N_CUSTOMERS).alias("customer_idx")
)
orders_with_fk = orders_df.join(customer_lookup, on="customer_idx")
undefinedorders_df = spark.range(N_ORDERS).select(
(F.abs(F.hash(F.col("id"))) % N_CUSTOMERS).alias("customer_idx")
)
orders_with_fk = orders_df.join(customer_lookup, on="customer_idx")
undefinedSetup
环境搭建
Requires Python 3.12 and databricks-connect>=16.4. Use :
uvbash
uv pip install "databricks-connect>=16.4,<17.4" faker numpy pandas holidays需要Python 3.12和databricks-connect>=16.4。使用:
uvbash
uv pip install "databricks-connect>=16.4,<17.4" faker numpy pandas holidaysRelated Skills
相关技能
- databricks-unity-catalog — Managing catalogs, schemas, and volumes
- databricks-dabs — DABs for production deployment
- databricks-unity-catalog — 管理Catalog、schema和卷
- databricks-dabs — 用于生产部署的DABs
Common Issues
常见问题
| Issue | Solution |
|---|---|
| Install locally: |
| Faker UDF is slow | Use |
| Out of memory | Increase |
| Referential integrity errors | Write master table to Delta first, read back for FK joins |
| NEVER use |
| Use |
| Broadcast variables not supported | NEVER use |
See references/2-troubleshooting.md for full troubleshooting guide.
| 问题 | 解决方案 |
|---|---|
| 本地安装: |
| Faker UDF运行缓慢 | 使用 |
| 内存不足 | 增加 |
| 引用完整性错误 | 先将主表写入Delta,再读取回来进行外键关联 |
| 无服务器环境绝不要使用 |
| 对于 |
| 广播变量不支持 | 无服务器环境绝不要使用 |
完整故障排除指南请参见references/2-troubleshooting.md。