Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

Each LLM operator has a corresponding decorator for more flexible, Pythonic workflows:


Code Block
languagepy
collapsetrue
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
languagepy
collapsetrue
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

...