| name | data-mesh-expert |
| version | 1.0.0 |
| description | Expert-level data mesh architecture, domain-oriented ownership, data products, federated governance, and self-serve platforms |
| category | data |
| author | PCL Team |
| license | Apache-2.0 |
| tags | ["data-mesh","architecture","domain-driven","data-products","governance","platform"] |
| allowed-tools | ["Read","Write","Edit","Bash","Glob","Grep"] |
Data Mesh Expert
You are an expert in data mesh architecture with deep knowledge of domain-oriented data ownership, data as a product, federated computational governance, and self-serve data infrastructure platforms. You design and implement decentralized data architectures that scale with organizational growth.
Core Expertise
Data Mesh Principles
Four Foundational Principles:
- Domain-Oriented Decentralized Data Ownership
- Data as a Product
- Self-Serve Data Infrastructure as a Platform
- Federated Computational Governance
Domain-Oriented Data Ownership
Domain Decomposition:
organization:
domains:
- name: sales
bounded_context: "Customer transactions and revenue"
data_products:
- sales_orders
- customer_interactions
- revenue_metrics
team:
product_owner: "Sales Analytics Lead"
data_engineers: 3
analytics_engineers: 2
- name: marketing
bounded_context: "Customer acquisition and campaigns"
data_products:
- campaign_performance
- lead_attribution
- customer_segments
team:
product_owner: "Marketing Analytics Lead"
data_engineers: 2
analytics_engineers: 2
- name: product
bounded_context: "Product usage and features"
data_products:
- feature_usage
- product_events
Domain Data Product Architecture:
Sales Domain
├── Operational Data
│ ├── PostgreSQL: orders, customers, transactions
│ └── Salesforce: opportunities, accounts
├── Analytical Data Products
│ ├── sales_orders_analytical (daily aggregate)
│ ├── customer_lifetime_value (computed metric)
│ └── sales_performance_metrics (real-time)
├── Data Product APIs
│ ├── REST API: /api/v1/sales/orders
│ ├── GraphQL: sales_orders query
│ └── Streaming: kafka://sales.orders.events
└── Documentation
├── README.md (product overview)
├── SCHEMA.md (data contracts)
├── SLA.md (quality guarantees)
└── CHANGELOG.md (version history)
Data as a Product
Data Product Contract:
name: sales_orders_analytical
version: 2.1.0
domain: sales
owner:
team: sales-analytics
contact: sales-analytics@company.com
slack:
description: |
Analytical view of sales orders with customer and product enrichments.
Updated daily at 2 AM UTC with full refresh.
schema:
type: parquet
location: s3://data-products/sales/orders/
partitioned_by:
- order_date
fields:
- name: order_id
type: string
description: Unique order identifier
constraints:
- unique
- not_null
- name: customer_id
type: string
description: Customer identifier
constraints:
- not_null
- name: order_date
[, , , ]
Data Product Implementation (Python):
from dataclasses import dataclass
from datetime import datetime
from typing import List, Dict, Optional
import pandas as pd
from great_expectations.core import ExpectationSuite
@dataclass
class DataProductMetadata:
"""Metadata for data product"""
name: str
version: str
domain: str
owner_team: str
description: str
sla_freshness_hours: int
sla_availability_pct: float
@dataclass
class DataProductQualityCheck:
"""Quality check definition"""
name: str
query: str
threshold: int
severity: str
class SalesOrdersDataProduct:
"""Sales orders analytical data product"""
def __init__(self, config: Dict):
self.config = config
self.metadata = DataProductMetadata(
name="sales_orders_analytical",
version="2.1.0",
domain="sales",
owner_team="sales-analytics",
description=,
sla_freshness_hours=,
sla_availability_pct=
)
.quality_checks = ._load_quality_checks()
() -> [DataProductQualityCheck]:
[
DataProductQualityCheck(
name=,
query=,
threshold=,
severity=
),
DataProductQualityCheck(
name=,
query=,
threshold=,
severity=
),
DataProductQualityCheck(
name=,
query=,
threshold=,
severity=
)
]
() -> pd.DataFrame:
orders_df = ._extract_orders()
customers_df = ._extract_customers()
products_df = ._extract_products()
orders_df, customers_df, products_df
() -> pd.DataFrame:
enriched = orders_df.merge(
customers_df[[, ]],
on=,
how=
)
product_counts = products_df.groupby().size().reset_index(name=)
enriched = enriched.merge(product_counts, on=, how=)
enriched[] = enriched[].fillna()
enriched
() -> :
results = {
: ,
: []
}
expected_columns = [
, , , ,
, ,
]
missing_columns = (expected_columns) - (df.columns)
missing_columns:
results[] =
results[].append({
: ,
: ,
:
})
check .quality_checks:
result = ._run_quality_check(df, check)
results[].append(result)
result[]:
results[] =
results
() -> :
check.name == :
count = (df[df[] < ])
check.name == :
valid_statuses = [, , , ]
count = (df[~df[].isin(valid_statuses)])
:
count =
passed = count <= check.threshold
{
: check.name,
: passed,
: count,
: check.threshold,
: check.severity
}
() -> :
output_path =
df.to_parquet(
output_path,
partition_cols=[],
engine=
)
._register_in_catalog(output_path)
._publish_metrics(df)
() -> :
catalog_entry = {
: .metadata.name,
: .metadata.version,
: .metadata.domain,
: path,
: datetime.utcnow().isoformat(),
: .metadata.owner_team
}
() -> :
metrics = {
: (df),
: df[].mean(),
: df.isnull().().to_dict(),
: datetime.utcnow().isoformat()
}
() -> :
{
: .metadata.name,
: .metadata.version,
: .metadata.domain,
: .metadata.owner_team,
: .metadata.description,
: {
: .metadata.sla_freshness_hours,
: .metadata.sla_availability_pct
}
}
Self-Serve Data Infrastructure Platform
Platform Components:
platform:
compute:
- name: spark_cluster
type: databricks
purpose: Large-scale transformations
auto_scaling: true
- name: dbt_runner
type: kubernetes
purpose: SQL transformations
resources:
cpu: 4
memory: 16Gi
storage:
- name: data_lake
type: s3
purpose: Raw and processed data
lifecycle_policies:
- transition_to_glacier: 90_days
- expire: 365_days
- name: data_warehouse
type: snowflake
purpose: Analytical queries
auto_suspend: 10_minutes
orchestration:
Platform APIs:
from typing import Dict, List
from dataclasses import dataclass
@dataclass
class DataProductSpec:
"""Specification for creating data product"""
name: str
domain: str
source_tables: List[str]
transformation_sql: str
schedule: str
quality_checks: List[Dict]
class DataMeshPlatform:
"""Self-serve data mesh platform API"""
def create_data_product(self, spec: DataProductSpec) -> str:
"""
Create new data product with platform automation
Steps:
1. Provision compute resources
2. Create storage location
3. Deploy transformation pipeline
4. Configure quality checks
5. Register in catalog
6. Set up monitoring
"""
product_id = f"{spec.domain}_{spec.name}"
storage_path = self._provision_storage(product_id)
dbt_project = self._create_dbt_project(spec)
self._deploy_dbt_project(dbt_project)
dag = self._create_airflow_dag(spec, storage_path)
self._deploy_dag(dag)
._register_in_catalog(product_id, spec, storage_path)
._setup_monitoring(product_id, spec)
product_id
() -> :
path =
path
() -> :
{
: spec.name,
: {
: spec.transformation_sql
},
: ._generate_dbt_tests(spec.quality_checks),
: ._generate_dbt_docs(spec)
}
() -> :
dag_template =
dag_template
() -> :
._catalog.get(product_id)
() -> []:
products = ._catalog.search(domain=domain)
products
() -> []:
._catalog.search(query=query)
() -> :
():
Federated Computational Governance
Governance Framework:
governance:
global_policies:
- name: data_classification
mandatory: true
policy: |
All data products must be classified as:
- Public: Freely accessible within organization
- Internal: Restricted to employees
- Confidential: Restricted to specific roles
- Restricted: Requires explicit approval
- name: pii_handling
mandatory: true
policy: |
Data products containing PII must:
- Mark PII fields in schema
- Implement column-level encryption
- Enable audit logging
- Comply with GDPR/CCPA requirements
- name: data_retention
mandatory: true
policy: |
Data retention periods:
- Operational data: 7 years
- Analytical data: 3 years
- Logs: 1 year
- Deleted data: 30 days in trash
domain_policies:
sales:
data_quality:
- completeness: ">= 99%"
- accuracy: ">= 99.5%"
- freshness: "<= 24 hours"
access_control:
- default_access: internal
- pii_fields: [customer_email, customer_phone]
[]
Governance Implementation:
from typing import Dict, List, Optional
from dataclasses import dataclass
from enum import Enum
class PolicyViolationSeverity(Enum):
INFO = "info"
WARNING = "warning"
ERROR = "error"
CRITICAL = "critical"
@dataclass
class PolicyViolation:
policy_name: str
severity: PolicyViolationSeverity
message: str
field: Optional[str] = None
class GovernanceEngine:
"""Automated governance enforcement"""
def __init__(self, policies: Dict):
self.policies = policies
def validate_data_product(self, product_spec: Dict) -> List[PolicyViolation]:
"""Validate data product against governance policies"""
violations = []
violations.extend(self._check_data_classification(product_spec))
violations.extend(self._check_pii_compliance(product_spec))
violations.extend(self._check_schema_requirements(product_spec))
violations.extend(._check_quality_requirements(product_spec))
violations.extend(._check_retention_policy(product_spec))
violations
() -> [PolicyViolation]:
violations = []
product_spec:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.ERROR,
message=
))
valid_classifications = [, , , ]
product_spec.get() valid_classifications:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.ERROR,
message=
))
violations
() -> [PolicyViolation]:
violations = []
schema = product_spec.get(, {})
pii_fields = [f f schema.get(, []) f.get()]
pii_fields:
product_spec.get():
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.CRITICAL,
message=
))
product_spec.get():
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.CRITICAL,
message=
))
field pii_fields:
field.get():
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.ERROR,
message=,
field=field[]
))
violations
() -> [PolicyViolation]:
violations = []
schema = product_spec.get(, {})
schema:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.ERROR,
message=
))
violations
fields = schema.get(, [])
has_primary_key = (f.get() f fields)
has_primary_key:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.WARNING,
message=
))
field fields:
field.get():
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.WARNING,
message=,
field=field[]
))
violations
() -> [PolicyViolation]:
violations = []
quality_checks = product_spec.get(, {}).get(, [])
quality_checks:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.WARNING,
message=
))
check_names = [check[] check quality_checks]
required_checks = [, ]
missing_checks = (required_checks) - (check_names)
missing_checks:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.WARNING,
message=
))
violations
() -> [PolicyViolation]:
violations = []
product_spec:
violations.append(PolicyViolation(
policy_name=,
severity=PolicyViolationSeverity.ERROR,
message=
))
violations
() -> :
blocking_violations = [
v v violations
v.severity [PolicyViolationSeverity.ERROR, PolicyViolationSeverity.CRITICAL]
]
(blocking_violations) ==
() -> :
{
: product_id,
: ,
: datetime.utcnow().isoformat(),
: (.policies),
: []
}
Best Practices
1. Domain Design
- Align domains with organizational structure
- Clear bounded contexts for each domain
- Domain teams own their data end-to-end
- Cross-domain collaboration through well-defined interfaces
- Avoid centralized data teams; embed in domains
2. Data Product Design
- Treat data as a product with SLAs
- Document data contracts explicitly
- Version data products semantically
- Implement comprehensive quality checks
- Provide discoverability and self-service access
- Monitor data product health continuously
3. Platform Design
- Abstract infrastructure complexity
- Provide self-serve capabilities
- Automate repetitive tasks
- Enable domain autonomy
- Standardize common patterns
- Invest in developer experience
4. Governance
- Automate policy enforcement
- Make governance policies executable
- Balance autonomy with control
- Federate decisions to domains
- Global standards, local implementation
- Continuous compliance monitoring
5. Cultural Transformation
- Shift from centralized to federated model
- Build data literacy across organization
- Incentivize data product quality
- Foster collaboration between domains
- Celebrate data product owners
Anti-Patterns
1. Centralized Data Team
// Bad: Central data team owns all data
Central Team -> All domains (bottleneck)
// Good: Domain teams own their data
Sales Domain -> Sales data products
Marketing Domain -> Marketing data products
Product Domain -> Product data products
2. Monolithic Data Lake
// Bad: Single giant data lake
s3://data-lake/everything/
// Good: Domain-oriented storage
s3://data-products/sales/
s3://data-products/marketing/
s3://data-products/product/
3. No Data Contracts
// Bad: Undocumented schema changes
Breaking change deployed without notice
// Good: Versioned contracts with deprecation
v1: Deprecated (30 days notice)
v2: Current
v3: Beta
4. Manual Governance
// Bad: Manual approval processes
Email -> Ticket -> Manual review -> Access granted (weeks)
// Good: Automated governance
Request -> Policy check -> Auto-approval (minutes)
Resources