DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
Each LLM operator has a corresponding decorator for more flexible, Pythonic workflows:
| Code Block | ||||
|---|---|---|---|---|
| ||||
from airflow.providers.ai.decorators import task
# Schema comparison with custom logic
@task.llm_schema_compare(data_sources=[s3_asset, postgres_asset])
def intelligent_schema_validation():
# Pre-processing: check business calendar
is_migration_window = check_migration_window()
if is_migration_window:
prompt = "Compare schemas and generate migration plan for scheduled maintenance window"
else:
prompt = "Compare schemas and flag breaking changes - no migrations allowed"
return {
"prompt": prompt,
"migration_allowed": is_migration_window,
"additional_context": {"maintenance_window": is_migration_window}
}
# Data quality with dynamic rules
@task.llm_data_quality(data_sources=[customer_asset])
def adaptive_quality_checks():
# Pre-processing: get current business rules
current_rules = fetch_business_rules()
seasonal_adjustments = get_seasonal_data_patterns()
prompt = f"""
Validate customer data against current business rules:
{current_rules}
Apply seasonal adjustments for data volume expectations:
{seasonal_adjustments}
Generate appropriate validation queries.
"""
return {
"prompt": prompt,
"business_rules": current_rules,
"seasonal_context": seasonal_adjustments
}
# File analysis with preprocessing
@task.llm_file_analysis(data_sources=[log_files_asset])
def analyze_logs_with_context():
# Pre-processing: get system context
recent_deployments = get_recent_deployments()
system_alerts = get_active_alerts()
prompt = f"""
Analyze log files for anomalies, considering:
- Recent deployments: {recent_deployments}
- Active system alerts: {system_alerts}
Focus on correlation between deployment events and error patterns.
"""
return {
"prompt": prompt,
"deployment_context": recent_deployments,
"alert_context": system_alerts
} |
...
| Code Block |
|---|
# Automatically injected context for PostgreSQL:
{
"database_type": "postgresql",
"version": "15.2",
"available_tables": ["customers", "orders", "products"],
"schema_info": {
"customers": {
"customer_id": {"type": "integer", "nullable": False, "primary_key": True},
"email_address": {"type": "varchar(255)", "nullable": False, "unique": True},
"total_revenue": {"type": "decimal(10,2)", "nullable": True}
}
},
"sample_data": {
"customers": [
{"customer_id": 1, "email_address": "john@example.com", "total_revenue": 1250.00},
{"customer_id": 2, "email_address": "jane@example.com", "total_revenue": 890.50}
]
},
"dialect_features": {
"supports_window_functions": True,
"supports_cte": True,
"date_functions": ["DATE_TRUNC", "EXTRACT", "AGE"]
}
}
File Operator Context (via S3Hook/GCSHook integration):
# Automatically injected context for S3 Parquet files:
{
"storage_type": "s3",
"file_format": "parquet",
"file_size_mb": 245,
"estimated_rows": 1000000,
"schema_info": {
"id": "int64",
"name": "string",
"email": "string",
"signup_date": "timestamp[ns]"
},
"sample_data": [
{"id": 1, "name": "John Doe", "email": "john@example.com"},
{"id": 2, "name": "Jane Smith", "email": "jane@example.com"}
],
"partitioning": ["year", "month"],
"compression": "snappy"
} |
...
We propose both embedded and separate HITL patterns:
Option A: Embedded HITL
| Code Block |
|---|
quality_check = LLMDataQualityOperator( |
...
task_id="customer_quality_analysis", |
...
data_sources=[customer_s3], |
...
prompt="Generate data quality validation queries", |
...
require_approval=True, # Built-in HITL |
...
approval_timeout=timedelta(hours=2) |
...
) |
Option B: Separate HITL Steps
| Code Block |
|---|
# Generate queries |
...
generate_queries = LLMDataQualityOperator( |
...
task_id="generate_quality_queries", |
...
data_sources=[customer_s3], |
...
prompt="Generate data quality validation queries", |
...
dry_run=True # Don't execute, just generate |
...
) |
...
# Human approval step |
...
approve_queries = ApprovalOperator( |
...
task_id="approve_queries", |
...
body="{{ ti.xcom_pull(task_ids='generate_quality_queries') }}", |
...
allow_modifications=True # Users can edit generated queries |
...
)
# Execute approved queries
) # Execute approved queries execute_analysis = AnalyticsOperator( |
...
task_id="execute_quality_checks", |
...
query="{{ ti.xcom_pull(task_ids='approve_queries') }}", |
...
data_sources=[customer_s3] |
...
)
) generate_queries >> approve_queries >> execute_analysis |
Complete Workflow Example
Here's a real-world scenario combining all components:
from datetime import datetime, timedelta
...
| Code Block | ||||
|---|---|---|---|---|
| ||||
from datetime import datetime, timedelta from airflow.sdk import DAG, Asset |
...
from airflow.providers.ai.operators import ( |
...
LLMSchemaCompareOperator, |
...
LLMDataQualityOperator, |
...
AnalyticsOperator |
...
) |
...
from airflow.operators.approval import ApprovalOperator |
...
# Define multi-cloud assets with structured metadata |
...
customer_s3 = Asset( |
...
name="customer_feed_s3", |
...
uri="s3://data-lake/customer/", |
...
conn_id="aws_default", |
...
schema={"id": "int32", "name": "string", "email": "string"}, |
...
sensitivity="pii", |
...
format="parquet" |
...
)
) customer_postgres = Asset( |
...
name="customer_master_postgres", |
...
uri="postgres://warehouse/public/customers", |
...
conn_id="postgres_default", |
...
schema={"customer_id": "integer", "full_name": "varchar", "email_address": "varchar"}, |
...
sensitivity="pii" |
...
)
...
) with DAG( |
...
"intelligent_data_validation", |
...
start_date=datetime(2024, 1, 1), |
...
schedule=timedelta(hours=6), |
...
) as dag: |
...
# 1. Detect schema drift between S3 feed and PostgreSQL master |
...
schema_drift = LLMSchemaCompareOperator( |
...
task_id="detect_schema_drift", |
...
data_sources=[customer_s3, customer_postgres], |
...
prompt="Identify schema mismatches that would break data loading", |
...
output_format="structured_report" |
...
) |
...
# 2. Generate data quality queries for new S3 data |
...
generate_quality_checks = LLMDataQualityOperator( |
...
task_id="generate_quality_queries", |
...
data_sources=[customer_s3], |
...
prompts=[ |
...
"Generate summary statistics queries", |
...
"Check for duplicate email addresses", |
...
"Validate data against business rules"
], |
...
dry_run=True |
...
) |
...
...
# 3. Human approval for generated queries (with edit capability) |
...
approve_queries = ApprovalOperator( |
...
task_id="approve_quality_queries", |
...
body="{{ ti.xcom_pull(task_ids='generate_quality_queries') }}", |
...
allow_modifications=True, |
...
timeout=timedelta(hours=2) |
...
) |
...
...
# 4. Execute approved quality checks using DataFusion |
...
execute_quality_checks = AnalyticsOperator( |
...
task_id="run_quality_analysis", |
...
query="{{ ti.xcom_pull(task_ids='approve_quality_queries') }}", |
...
data_sources=[customer_s3], |
...
engine="datafusion", |
...
output_location="s3://results/quality-reports/" |
...
) |
...
schema_drift >> generate_quality_checks >> approve_queries >> execute_quality_checks |
Evolution Path: From LLMOperator to AITask
...
- Dynamic prompt generation based on runtime conditions
- Custom pre-processing logic before LLM calls
- Context enrichment from external systems (business rules, calendars, alerts)
- Conditional logic for different operational scenarios
Enhanced Asset System
# Proposed Asset structure evolution
...
extra: Dict[str, Any] = field(default_factory=dict) # Backward compatibility
AnalyticsOperator with DataFusion
- Unified interface for multi-cloud data access across S3, GCS, Azure Blob Storage
- Native handling of Parquet, Iceberg, Delta Lake, JSON, CSV, Avro formats
- High-performance processing: 50M+ records in 10-15 seconds on single node(This is from my experiments)
- Cost-effective alternative to Spark for AI-generated analytical queries
- SQL dialect unification - write once, run anywhere
- DataFusion table provider supports integrating existing databases like sqlite, postgres. A good option here is using datafusion-table-providers. https://github.com/datafusion-contrib/datafusion-table-providers
Flexible HITL Integration
- Embedded approval within operators (require_approval=True)
- Separate ApprovalOperator with query modification capabilities
- Configurable timeout and escalation policies
...