id: ai-synthetic-relational-data-generator
namespace: company.data
description: |
Generates privacy-preserving synthetic relational datasets without PII leakage,
validates foreign-key referential integrity and temporal consistency in DuckDB,
and publishes QA-ready artifacts.
triggers:
- id: scheduled_refresh
type: io.kestra.plugin.core.trigger.Schedule
description: Weekly synthetic data generation sweep for staging environments.
Shipped disabled by default.
cron: "0 6 * * 1"
disabled: true
- id: webhook_pipeline_trigger
type: io.kestra.plugin.core.trigger.Webhook
description: Authenticated webhook for CI/CD staging seeding pipelines.
key: "{{ secret('WEBHOOK_KEY') }}"
inputs:
- id: user_count
type: INT
defaults: 25
description: Number of synthetic user profiles to generate.
- id: max_orders_per_user
type: INT
defaults: 4
description: Maximum number of synthetic orders per user.
tasks:
- id: generate_relational_entities
type: io.kestra.plugin.scripts.python.Script
description: Generates realistic multi-table relational datasets with strictly
synthetic attributes and zero PII.
taskRunner:
type: io.kestra.plugin.core.runner.Process
script: |
import random
import csv
import uuid
from datetime import datetime, timedelta
user_count = int("{{ inputs.user_count }}")
max_orders = int("{{ inputs.max_orders_per_user }}")
now = datetime.now()
users = []
orders = []
payments = []
tiers = ["STARTER", "PROFESSIONAL", "ENTERPRISE"]
methods = ["CREDIT_CARD", "ACH_TRANSFER", "CORPORATE_INVOICE"]
for i in range(user_count):
uid = f"usr_{uuid.uuid4().hex[:10]}"
created_at = (now - timedelta(days=random.randint(15, 60))).isoformat()
users.append([uid, f"tenant-{i+1}.internal", random.choice(tiers), created_at])
# Generate related orders
order_count = random.randint(1, max_orders)
for _ in range(order_count):
oid = f"ord_{uuid.uuid4().hex[:10]}"
order_time = (now - timedelta(days=random.randint(1, 14))).isoformat()
amount = round(random.uniform(75.0, 1250.0), 2)
orders.append([oid, uid, amount, "COMPLETED", order_time])
# Generate related payment
pid = f"pay_{uuid.uuid4().hex[:10]}"
payments.append([pid, oid, random.choice(methods), amount, order_time, "SETTLED"])
# Write synthetic CSV tables
with open("synthetic_users.csv", "w", newline="") as f:
w = csv.writer(f)
w.writerow(["user_id", "company_domain", "account_tier", "created_at"])
w.writerows(users)
with open("synthetic_orders.csv", "w", newline="") as f:
w = csv.writer(f)
w.writerow(["order_id", "user_id", "amount", "status", "order_timestamp"])
w.writerows(orders)
with open("synthetic_payments.csv", "w", newline="") as f:
w = csv.writer(f)
w.writerow(["payment_id", "order_id", "payment_method", "amount", "payment_timestamp", "settlement_status"])
w.writerows(payments)
print(f"Generated {len(users)} users, {len(orders)} orders, and {len(payments)} payments with referential linkage.")
outputFiles:
- synthetic_users.csv
- synthetic_orders.csv
- synthetic_payments.csv
- id: validate_referential_integrity
type: io.kestra.plugin.jdbc.duckdb.Query
description: Executes in-memory DuckDB query verifying foreign-key constraints
and temporal consistency across generated tables.
communityExtensions: []
inputFiles:
synthetic_users.csv: "{{
outputs.generate_relational_entities.outputFiles['synthetic_users.csv']
}}"
synthetic_orders.csv: "{{
outputs.generate_relational_entities.outputFiles['synthetic_orders.csv']
}}"
synthetic_payments.csv: "{{
outputs.generate_relational_entities.outputFiles['synthetic_payments.cs\
v'] }}"
fetchType: FETCH_ONE
sql: |
WITH users AS (
SELECT * FROM read_csv_auto('synthetic_users.csv', header = true)
),
orders AS (
SELECT * FROM read_csv_auto('synthetic_orders.csv', header = true)
),
payments AS (
SELECT * FROM read_csv_auto('synthetic_payments.csv', header = true)
),
order_fk_audit AS (
SELECT
COUNT(CASE WHEN u.user_id IS NULL THEN 1 END) AS missing_user_fks
FROM orders o
LEFT JOIN users u ON o.user_id = u.user_id
),
payment_fk_audit AS (
SELECT
COUNT(CASE WHEN o.order_id IS NULL THEN 1 END) AS missing_order_fks
FROM payments p
LEFT JOIN orders o ON p.order_id = o.order_id
),
entity_counts AS (
SELECT
(SELECT COUNT(*) FROM users) AS user_count,
(SELECT COUNT(*) FROM orders) AS order_count,
(SELECT COUNT(*) FROM payments) AS payment_count
)
SELECT
c.user_count,
c.order_count,
c.payment_count,
o.missing_user_fks,
p.missing_order_fks,
(o.missing_user_fks + p.missing_order_fks) AS total_fk_violations,
CASE
WHEN (o.missing_user_fks + p.missing_order_fks) = 0 THEN true
ELSE false
END AS is_dataset_valid
FROM entity_counts c
CROSS JOIN order_fk_audit o
CROSS JOIN payment_fk_audit p;
- id: data_quality_gate
type: io.kestra.plugin.core.flow.If
description: Halts staging deployment if synthetic datasets fail foreign-key
constraints.
condition: "{{ outputs.validate_referential_integrity.row.is_dataset_valid == true }}"
then:
- id: log_validation_pass
type: io.kestra.plugin.core.log.Log
description: Logs dataset verification summary.
message: "Synthetic dataset verified: {{
outputs.validate_referential_integrity.row.user_count }} users, {{
outputs.validate_referential_integrity.row.order_count }} orders, {{
outputs.validate_referential_integrity.row.payment_count }} payments.
Zero foreign-key violations. Ready for QA staging."
- id: notify_data_team_success
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts Data Engineering Slack channel of verified synthetic dataset
readiness.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#data-quality-alerts",
"text": "✅ *Synthetic Relational Dataset Ready*: Generated {{ outputs.validate_referential_integrity.row.user_count }} users, {{ outputs.validate_referential_integrity.row.order_count }} orders, and {{ outputs.validate_referential_integrity.row.payment_count }} payments with zero constraint violations. Artifacts archived."
}
else:
- id: alert_integrity_violation
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts Data Engineering Slack channel if referential integrity
check fails.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#data-quality-alerts",
"text": "🚨 *Synthetic Data Quality Gate Failed*: Discovered {{ outputs.validate_referential_integrity.row.total_fk_violations }} foreign-key constraint violations in generated dataset. Staging push aborted."
}
errors:
- id: alert_on_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts on-call engineering channel if data generation or DuckDB
validation fails unexpectedly.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#data-ops",
"text": "⚠️ *Pipeline Error*: Synthetic data generation flow `{{ flow.id }}` failed on execution `{{ execution.id }}`."
}