Skip to content

Airflow Data Quality Provider part 1#69575

Open
gopidesupavan wants to merge 19 commits into
apache:mainfrom
gopidesupavan:dq-provider-backend
Open

Airflow Data Quality Provider part 1#69575
gopidesupavan wants to merge 19 commits into
apache:mainfrom
gopidesupavan:dq-provider-backend

Conversation

@gopidesupavan

@gopidesupavan gopidesupavan commented Jul 7, 2026

Copy link
Copy Markdown
Member

Adds a new apache-airflow-providers-common-dataquality provider for DbApiHook-based data quality checks.

Airflow already has SQL check operators, and many users rely on them for data quality today. This provider adds a DQRule / RuleSet layer for checks that need stable rule identity, persisted history, and a connection to Airflow assets. That makes quality results easier to analyze over time, lets downstream asset consumers gate on recent quality, and gives LLM-assisted workflows one schema to generate when proposing checks from table context. Execution still goes through existing common.sql / DbApiHook connections.

This PR is the backend/provider slice only. The UI plugin and read-only API are intentionally left for a follow-up PR.

Ships:

  • DQRule and RuleSet models for named data quality rules.
  • Built-in SQL checks for common table and column checks, executed through common.sql / DbApiHook.
  • custom_sql support for database-specific or more complex checks.
  • DQCheckOperator and the @task.dq_check TaskFlow decorator.
  • A configurable results backend under [dq] results_path for task, run, and rule-level history.
  • Experimental asset helpers, asset_quality() and require_quality(), that attach provider-owned quality metadata to assets without changing Airflow core.
  • Documentation and example Dags covering end-to-end usage with and without LLM-generated rules.
  • A DQ rule-authoring skill that LLM-assisted workflows can use to generate rules from table/schema context:
    https://github.com/gopidesupavan/airflow/blob/f32940bd261b94238256eaced9150dd51329ce3e/providers/dq/src/airflow/providers/dq/skills/dq-rule-authoring/SKILL.md

This first version is intentionally focused on the backend contract: deterministic rule definitions, SQL execution through existing Airflow SQL providers, persisted results, and asset-linked quality summaries.

Design decisions:

  • Results are stored through an object-storage/local-file backend instead of adding new metadata DB tables in the first provider drop. This keeps the provider self-contained, avoids Airflow core migrations, and lets deployments choose a durable store such as S3, GCS, or local files via [dq] results_path.
  • The backend stores keyed JSON records for task runs, task instances, and per-rule history so later readers, including a future UI/API layer, can access common views without scanning unrelated runs.
  • Asset support is implemented with provider-owned metadata, not Airflow core changes. Static quality configuration is attached to Asset.extra["airflow.dq"]; runtime summaries are attached to asset events under extra["airflow.dq.result"].
  • The first release starts with DbApiHook / SQL execution because Airflow already has broad database coverage through common.sql. File and object-store data checks are left for a later iteration.

Later iterations:

  • Read-only API and minimal Airflow UI plugin for viewing task/run results and rule history.
  • File/object-store based checks, where Airflow reads data from S3/GCS/local files or other object stores and runs quality rules directly against that data.
  • OpenLineage integration for data quality facets.
  • More built-in checks.
  • Trigger support for DQ checks.
  • DQProfileOperator.

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)
    codex

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@gopidesupavan gopidesupavan changed the title Airflow Data Quality Provider Airflow Data Quality Provider Jul 7, 2026
@gopidesupavan gopidesupavan changed the title Airflow Data Quality Provider Airflow Data Quality Provider part 1 Jul 7, 2026
@gopidesupavan
gopidesupavan requested a review from o-nikolas July 7, 2026 21:42
@gopidesupavan

gopidesupavan commented Jul 7, 2026

Copy link
Copy Markdown
Member Author

This is part 1 , the UI plugin is part of this big PR #69413, i will ship that separate. but to show how the UI looks like here is some screenshots #69413 (comment) (this is minimal what had done with help of UI not an UI expert anyone can modify or help with this 😄 )

@gopidesupavan
gopidesupavan force-pushed the dq-provider-backend branch 2 times, most recently from cbdf1fd to 9dac869 Compare July 8, 2026 07:24
@gopidesupavan
gopidesupavan requested review from kaxil and vikramkoka July 8, 2026 09:40
@gopidesupavan
gopidesupavan force-pushed the dq-provider-backend branch 3 times, most recently from ef557f3 to 106b8de Compare July 10, 2026 21:25
@gopidesupavan
gopidesupavan requested a review from shahar1 July 10, 2026 21:25
@gopidesupavan
gopidesupavan force-pushed the dq-provider-backend branch 3 times, most recently from 377ca0a to 1047525 Compare July 12, 2026 19:38
@gopidesupavan
gopidesupavan requested review from Lee-W and removed request for Lee-W July 16, 2026 07:43
@gopidesupavan

Copy link
Copy Markdown
Member Author

sorry @Lee-W i meant to add and it got removed, readded back.

Comment thread providers/common/dataquality/docs/rules.rst Outdated
config["conn_id"] = conn_id
if table:
config["table"] = table
asset.extra[DQ_EXTRA_KEY] = config

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please don't use asset.extra. This is going to be removed #55200.

@gopidesupavan

Copy link
Copy Markdown
Member Author

@Lee-W regarding extra what other options you would suggest? do we have any other ways to store that now?

@gopidesupavan
gopidesupavan requested a review from Dev-iL July 16, 2026 20:12
@gopidesupavan
gopidesupavan requested review from Lee-W and kaxil July 17, 2026 11:02

@SameerMesiah97 SameerMesiah97 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since this is a very large PR, I could only review a few pieces of it but there is some general feedback I can share:

  1. In your last dev list response, you mentioned that the provider currently builds on DbApiHook, but that DataFusion is the next execution engine you plan to add. That made me wonder about the long-term architecture. Since DataFusion isn't a DB-API implementation, I would expect the provider to evolve towards a more general execution layer that both SQL and DataFusion plug into. I am just wondering what this provider would like when this happens?

  2. I think the docs could be made more objective and instructional. In several places the wording feels a little clunky or promotional, whereas I'd expect Airflow documentation to focus on clearly explaining how to use the feature. Simplifying some of the phrasing would also improve readability. This may come across as a bit pedantic but we have to keep in mind that docs will be the first point of interaction for this provider. LLMs can be very useful for this with good prompts if you think it will be too time-consuming.

  3. I spotted some tests which seemed redundant, effectively testing native python behaviour. I would review the other tests to see if they can be trimmed.

Overall, I think this is a very interesting addition at the provider level. Not your typical external service integration.

.. exampleinclude:: /../src/airflow/providers/common/dataquality/example_dags/example_dq_llm_generated_ruleset.py
:language: python
:start-after: [START howto_decorator_dq_check_llm_runtime_ruleset]
:end-before: [END howto_decorator_dq_check_llm_runtime_ruleset]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I understand that this may come across as a bit pedantic but I think the documentation here veers too far into marketing language when I believe they should be purely instructional. Statements like 'Writing a RuleSet by hand for every table doesn't scale' is of course very defensible but it is a bit too opinionated for provider documentation.

I think it would not hurt to use an LLM here with specific instructions to keep it factual. I also found some of the phrasing a little awkward from a native English perspective. Again, a well-prompted LLM should be very useful for this sort of thing.

This feedback applies to the rest of the docs too. I glanced at them and they had similar issues.

:param asset: An asset decorated with :func:`~airflow.providers.common.dataquality.assets.asset_quality`.
Supplies defaults for ``ruleset``, ``table``, and ``conn_id`` (explicit arguments
win) and is automatically added to the task's outlets so its asset events carry
the check summary.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Now, this is a more general concern about the API of this operator: it is a bit hard to understand with the inter-relationships between the different parameters. For example, ruleset, and `table`` are all optional, but become required unless asset supplies them. I think we should construct the API in a way that removes all these interdependencies so that our users would find the operator more intuitive.

Maybe we should simplify the public API by making asset the sole source of configuration when it is provided, rather than allowing table, ruleset, and conn_id to override it. At the moment the operator supports multiple overlapping configuration mechanisms, which introduces precedence rules and conditional required parameters.

warned = [r.rule_name for r in results if r.status == WARN]
if self.fail_on == "error" and failed:
raise DQCheckFailedError(f"Data quality rules failed: {failed} (score={summary['score']})")
if self.fail_on == "warn" and (failed or warned):

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I am just curious why you are using raw strings such as "error" and "warn" here instead of an Enum? I see you added an Enum class in this PR. Perhaps, you could add a new class like this:

class FailOn(str, Enum):
    ERROR = "error"
    WARN = "warn"
    NEVER = "never"

Also, this more of a nit but I think we could handle each case more explicity like the below:

if self.fail_on is FailOn.ERROR:
    ...
elif self.fail_on is FailOn.WARN:
    ...
elif self.fail_on is FailOn.NEVER:
    self.log.warning(...)



def test_dq_check_failed_error_is_runtime_error():
assert isinstance(DQCheckFailedError("failed check"), RuntimeError)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are these tests needed? It seems like you are just testing inheritance.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants