Forum Discussion

codeautomation's avatar
3 days ago

How should Microsoft Fabric data pipelines handle AI-driven validation ?

I’m exploring patterns for adding AI-assisted validation into Microsoft Fabric data engineering workflows before transformed data is written to a Lakehouse.

For example:

Source → Fabric Pipeline / Notebook → Transformation → AI Validation → Lakehouse → Semantic Model

I’m particularly interested in:

  • validating schema changes before writes
  • detecting unusual or inconsistent values
  • deciding when to fail the pipeline vs continue with warnings
  • logging validation results for later review
  • retrying failed validation steps
  • keeping the process efficient for larger datasets

Would you normally implement this through Fabric notebooks, pipeline activities, or a separate validation service?

3 Replies

  • Hi codeautomation,

    The architecture above is the right shape, so I'll add a concrete notebook version of it. It covers your points in one place: schema checks before the write, fail vs warn, logging, retries and keeping it cheap on large data.

    1. Contract check before anything else, so schema drift fails fast with a clear message instead of surfacing later in the semantic model:

    from pyspark.sql import functions as F
    import json, uuid
    
    run_id = str(uuid.uuid4())
    df = spark.read.table("bronze.orders")
    
    expected = {"order_id": "int", "customer_id": "int",
                "amount": "decimal(18,2)", "order_date": "date"}
    actual = dict(df.dtypes)
    drift = {c: (t, actual.get(c)) for c, t in expected.items()
             if actual.get(c) != t}
    
    if drift:
        notebookutils.notebook.exit(json.dumps(
            {"status": "FAIL", "reason": "schema",
             "detail": str(drift), "run_id": run_id}))

    2. Deterministic rules with a severity each, evaluated in one pass with Spark:

    rules = {
        "null_key":    ("critical", F.col("order_id").isNull()),
        "neg_amount":  ("warning",  F.col("amount") < 0),
        "future_date": ("warning",  F.col("order_date") > F.current_date()),
    }
    
    flagged = df
    for name, (_, cond) in rules.items():
        flagged = flagged.withColumn(
            name, F.coalesce(cond, F.lit(False)))
    flagged = flagged.cache()
    
    counts = flagged.agg(*[
        F.sum(F.col(n).cast("int")).alias(n) for n in rules
    ]).first().asDict()
    
    dup_keys = (df.groupBy("order_id").count()
                  .filter("count > 1").count())
    
    critical = (counts["null_key"] or 0) + dup_keys
    warnings = sum(counts[n] or 0 for n, (s, _) in rules.items()
                   if s == "warning")

    3. Log every rule, every run, then quarantine instead of silently fixing:

    log = [(run_id, "orders", n, s, int(counts[n] or 0))
           for n, (s, _) in rules.items()]
    log.append((run_id, "orders", "duplicate_key", "critical", dup_keys))
    
    (spark.createDataFrame(
        log, "run_id string, tbl string, rule string, "
             "severity string, failed_rows long")
      .withColumn("run_ts", F.current_timestamp())
      .write.mode("append").saveAsTable("dq.validation_log"))
    
    any_issue = " OR ".join(rules)
    (flagged.filter(any_issue)
      .withColumn("run_id", F.lit(run_id))
      .write.mode("append").saveAsTable("dq.orders_quarantine"))
    
    if critical == 0:
        (flagged.filter(f"NOT ({any_issue})").drop(*rules)
          .write.mode("append").saveAsTable("silver.orders"))

    4. AI only on what the rules could not decide, and only once per row. Hashing the row content means retries and reruns never pay twice for the same record:

    suspects = (flagged.filter("neg_amount OR future_date")
      .withColumn("row_hash",
          F.sha2(F.to_json(F.struct(*df.columns)), 256)))
    
    reviewed = spark.read.table("dq.ai_review_cache").select("row_hash")
    to_review = suspects.join(reviewed, "row_hash", "left_anti")
    
    # run Fabric AI functions (classify / explain) on to_review only,
    # then append row_hash + verdict to dq.ai_review_cache

    5. Return one result to the pipeline:

    status = "FAIL" if critical else ("WARN" if warnings else "PASS")
    notebookutils.notebook.exit(json.dumps({
        "status": status, "run_id": run_id,
        "critical": critical, "warnings": warnings}))

    In the pipeline, the If Condition reads the exit value:

    @equals(
      json(activity('Validate').output.result.exitValue).status,
      'FAIL')

    For retries, I would set a retry policy on the notebook activity for technical errors (timeouts, throttling) only. A data-quality FAIL should never be retried, because the same data fails the same way.

    So my overall approach would be:

    contract check -> fail fast on schema drift
    -> Spark rules with severity
    -> log every rule, quarantine failed rows
    -> AI only on unresolved rows, cached by hash
    -> exit value drives pipeline branching

    The dq.validation_log table also gives you a simple Power BI view of data quality trends per table and rule over time.

    AI-assisted drafting: AI was used to help structure and phrase this response. I reviewed and validated the technical content before posting.

  • ShivekMaharaj's avatar
    ShivekMaharaj
    Icon for Community Champion rankCommunity Champion

    Hi codeautomation​,

    I would separate deterministic data-quality rules from AI-assisted validation rather than treating AI as the primary gate.

    For things such as:

    • schema changes
    • null checks
    • duplicate detection
    • allowed values
    • referential integrity
    • numeric/date ranges

    I would keep the validation deterministic.

    Then I would use AI for cases where the rule is harder to express precisely, for example:

    • classifying unusual records
    • detecting semantic inconsistencies
    • validating free-text fields
    • explaining anomalies
    • extracting information from unstructured data


    Fabric now has AI Functions that can be used from notebooks, Dataflow Gen2 and SQL for operations such as classification, extraction and similarity.

    Architecturally, I would normally use:

    Pipeline
    -> Bronze/raw data
    -> Notebook or validation step
    -> deterministic checks
    -> optional AI validation
    -> write validation results
    -> If Condition
    -> Silver/curated data or quarantine

    That also fits Microsoft’s medallion architecture guidance, where raw data is preserved in Bronze and validation/cleansing happens before data progresses into Silver.

    For orchestration, a notebook can return a validation result to the pipeline, and the pipeline can branch using an If Condition. Fabric supports notebook exit values specifically for this kind of conditional orchestration.

    I would also log the validation outcome rather than only returning pass/fail. For example:

    • run ID
    • rule name
    • severity
    • failed row count
    • sample failures
    • AI confidence where applicable
    • timestamp


    Then use severity levels such as:

    Critical
    -> fail pipeline / quarantine
    
    Warning
    -> continue but log
    
    Informational
    -> continue

    For scale, I would avoid sending every row through an LLM. Run deterministic Spark checks across the full dataset first, then send only suspicious or ambiguous records to the AI layer.

    That keeps the process cheaper, faster and easier to audit.

    So my preference would be:

    • Fabric Pipeline for orchestration
    • Fabric Notebook/Spark for large-scale validation
    • AI Functions only where semantic reasoning actually adds value
    • Lakehouse tables for validation logs and quarantined records

     

    AI-assisted drafting: AI was used to help structure and phrase this response. I reviewed and validated the technical content before posting.

  • v-achippa's avatar
    v-achippa
    Icon for Community Support rankCommunity Support

    Hi codeautomation​,

    Thank you for reaching out to Microsoft Fabric Community.

    Use the Fabric pipeline to orchestrate the process and a Fabric notebook to perform the validation. The notebook can check schema, nulls, duplicates, value ranges and anomalies before writing to the curated Lakehouse table.

    Return the validation status to the pipeline and use an If Condition activity. Critical validation failures can trigger a Fail activity, while noncritical issues can be logged as warnings and written to a quarantine table. Retry only temporary technical failures, not data-quality failures. For large datasets, validate only new or changed partitions and use Spark-based checks rather than processing records individually.

    Thanks and regards,
    Anjan Kumar Chippa