databricks-synthetic-data-gen

Compare original and translation side by side

🇺🇸

Original

English
🇨🇳

Translation

Chinese
Catalog 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

参考资料

WhenGuide
User mentions ML model training or complex time patternsreferences/1-data-patterns.md — ML-ready data, time multipliers, row coherence
Errors during generationreferences/2-troubleshooting.md — Fixing common issues
场景指南
用户提及ML模型训练或复杂时间模式references/1-data-patterns.md — 适用于ML的数据、时间乘数、行一致性
生成过程中出现错误references/2-troubleshooting.md — 修复常见问题

Critical Rules

关键规则

  1. Data tells a story — Something goes wrong, impacts $, can be analyzed and fixed. Show Databricks value.
  2. 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.
  3. Industry terms, simple schema — Use domain-specific vocabulary but keep it easy to understand (few tables, clear relationships)
  4. Never uniform distributions — Skewed categories, log-normal amounts, 80/20 patterns. Flat = no story = useless
  5. Enough data for trends — ~100K+ rows for main tables so patterns survive aggregation
  6. Ask for catalog/schema — Never default, always confirm before generating
  7. Present plan for approval — Show tables, distributions, assumptions before writing code
  8. Master tables first — Generate parent tables, write to Delta, then create children with valid FKs
  9. Use Spark + Faker + Pandas UDFs — Scalable, parallel. Polars only if user explicitly wants local + <30K rows
  10. 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)
  11. No
    .cache()
    or
    .persist()
    — Not supported on serverless. Write to Delta, read back for joins
  12. No Python loops or
    .collect()
    — Use Spark parallelism. No driver-side iteration, avoid Pandas↔Spark conversions
  1. 数据承载场景——出现问题,影响金额,可分析并修复。展示Databricks的价值。
  2. 所有数据服务于场景——每张表和每列都必须连贯,可用于仪表板或ML模型。无孤立数据,无随机噪声——如果无法帮助解释、绘制未来仪表板或进行预测,则不要生成。
  3. 行业术语,简洁schema——使用领域特定词汇,但保持易懂(少量表,关系清晰)
  4. 绝不使用均匀分布——偏斜分类、对数正态金额、80/20模式。均匀分布=无场景=无用
  5. 足够数据以呈现趋势——主表需约10万+行,确保聚合后模式仍存在
  6. 询问Catalog/Schema——绝不默认,生成前务必确认
  7. 提交计划供批准——编写代码前展示表、分布、假设
  8. 先生成主表——生成父表,写入Delta,再创建带有有效外键的子表
  9. 使用Spark + Faker + Pandas UDFs——可扩展、并行化。仅当用户明确要求本地生成且行数<30K时使用Polars
  10. 默认使用Databricks Connect Serverless生成数据——必要时在Python 3.12上更新databricks-connect(除非指示不使用Databricks Connect,否则避免使用execute_code)
  11. 禁止使用
    .cache()
    .persist()
    ——无服务器环境不支持。写入Delta后再读取进行关联
  12. 禁止Python循环或
    .collect()
    ——使用Spark并行化。避免驱动端迭代,减少Pandas与Spark之间的转换

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.
TableDescriptionRowsKey Assumptions
customersCustomer profiles with tier, MRR10,000Enterprise 10% but 60% of revenue
ticketsSupport tickets with priority, resolution_time80,000Spike during outage, SLA breaches
incidentsSystem events (outages, deployments)50Payment outage mid-month
churn_eventsCustomer cancellations with reason500Spike after poor support experience
Business metrics:
  • customers.mrr
    — Revenue at risk ($)
  • tickets.resolution_hours
    — SLA performance
  • churn_events.lost_mrr
    — Churn impact ($)
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
    — 面临风险的收入(美元)
  • tickets.resolution_hours
    — SLA表现
  • 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
databricks experimental aitools tools query
to validate generated data (row counts, distributions, referential integrity). Query parquet files directly:
bash
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.
使用
databricks experimental aitools tools query
验证生成的数据(行数、分布、引用完整性)。直接查询parquet文件:
bash
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 pd
python
from databricks.connect import DatabricksSession
from pyspark.sql import functions as F
from pyspark.sql.types import StringType
import pandas as pd

Setup 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 systems
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")

**按规模设置分区:** `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-PatternWhy It's SlowDo This Instead
Python loops on driverSingle-threaded, no parallelismUse
spark.range()
+ Spark operations
.collect()
then iterate
Brings all data to driver memoryKeep data in Spark, use DataFrame ops
Pandas → Spark → PandasSerialization overhead, defeats distributionStay in Spark, use
pandas_udf
only for UDFs
Read/write temp filesUnnecessary I/OChain DataFrame transformations
Scalar UDFsRow-by-row processingUse
pandas_udf
for batch processing
Good pattern:
spark.range()
→ Spark transforms →
pandas_udf
for Faker → write directly
生成的脚本必须具备高性能。绝对禁止以下操作:
反模式为何缓慢替代方案
驱动端Python循环单线程,无并行性使用
spark.range()
+ Spark操作
.collect()
后迭代
将所有数据加载到驱动端内存数据保留在Spark中,使用DataFrame操作
Pandas → Spark → Pandas序列化开销,失去分布式优势保持在Spark环境中,仅在UDF中使用
pandas_udf
读写临时文件不必要的I/O链式DataFrame转换
标量UDF逐行处理使用
pandas_udf
进行批处理
推荐模式:
spark.range()
→ Spark转换 →
pandas_udf
调用Faker → 直接写入

Common 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
np.random.lognormal(mean, sigma)
— always positive, long tail:
  • Enterprise:
    lognormal(7.5, 0.8)
    → ~$1800 median
  • Pro:
    lognormal(5.5, 0.7)
    → ~$245 median
  • Free:
    lognormal(4.0, 0.6)
    → ~$55 median
使用
np.random.lognormal(mean, sigma)
— 始终为正,长尾分布:
  • 企业客户:
    lognormal(7.5, 0.8)
    → 中位数约$1800
  • Pro客户:
    lognormal(5.5, 0.7)
    → 中位数约$245
  • 免费客户:
    lognormal(4.0, 0.6)
    → 中位数约$55

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
.cache()
on serverless):
python
undefined
先将主表写入Delta,再读取回来进行外键关联(无服务器环境禁止使用
.cache()
):
python
undefined

1. 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")
undefined
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")
undefined

Setup

环境搭建

Requires Python 3.12 and databricks-connect>=16.4. Use
uv
:
bash
uv pip install "databricks-connect>=16.4,<17.4" faker numpy pandas holidays
需要Python 3.12和databricks-connect>=16.4。使用
uv
bash
uv pip install "databricks-connect>=16.4,<17.4" faker numpy pandas holidays

Related 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

常见问题

IssueSolution
ModuleNotFoundError: faker
Install locally:
uv pip install faker
, import inside UDF
Faker UDF is slowUse
pandas_udf
for batch processing
Out of memoryIncrease
numPartitions
in
spark.range()
Referential integrity errorsWrite master table to Delta first, read back for FK joins
PERSIST TABLE is not supported on serverless
NEVER use
.cache()
or
.persist()
with serverless
- write to Delta table first, then read back
F.window
vs
Window
confusion
Use
from pyspark.sql.window import Window
for
row_number()
,
rank()
, etc.
F.window
is for streaming only.
Broadcast variables not supportedNEVER use
spark.sparkContext.broadcast()
with serverless
See references/2-troubleshooting.md for full troubleshooting guide.
问题解决方案
ModuleNotFoundError: faker
本地安装:
uv pip install faker
,在UDF内部导入
Faker UDF运行缓慢使用
pandas_udf
进行批处理
内存不足增加
spark.range()
中的
numPartitions
引用完整性错误先将主表写入Delta,再读取回来进行外键关联
PERSIST TABLE is not supported on serverless
无服务器环境绝不要使用
.cache()
.persist()
- 先写入Delta表,再读取回来
F.window
Window
混淆
对于
row_number()
rank()
等,使用
from pyspark.sql.window import Window
F.window
仅用于流处理。
广播变量不支持无服务器环境绝不要使用
spark.sparkContext.broadcast()
完整故障排除指南请参见references/2-troubleshooting.md