Airflow Data Quality Provider Plugin#69413
Conversation
0432ee7 to
6acc4da
Compare
5c4aff2 to
f32940b
Compare
a8c6dfd to
e0f7ff9
Compare
o-nikolas
left a comment
There was a problem hiding this comment.
This first version is intentionally small.
I'm not so sure that it is 😅 This is an absolutely enormous PR, 11.8k lines added across 143 files. This is quite difficult to review. Is it possible to ship this in smaller pieces?
hehe it has both UI and backend shipped thats the reason it looks like big.. and some docs 😄 I am happy to ship this into two separate PR's one with UI and another backend provider specific.. |
|
@o-nikolas here is the part 1 #69575 hope its fine 😄 now.. as this is provider starting so it has autogenrated breeze docs and content in documentation .. |
5b30192 to
54ba129
Compare
|
Plugin UI is looking good. Would it be helpful to finally add plugins to the asset page to also show Data Quality there? |
Yeah thought of that it seems currently it doesnt support?, are you referring to the changes in core to support plugins in asset page ? |
I would open a PR to add support. |
thanks.. :) |
503448b to
0ccbe42
Compare
0ccbe42 to
8f2f582
Compare
Dev-iL
left a comment
There was a problem hiding this comment.
Left a few comments on the docs.
One suggestion: in terms of visualizing the results of a validation, a user might not always be interested in a strict pass/fail, but to be able to visualize "how bad the situation is", so let's say a user specifies a colormap or at least two colors and then then according to the validation score the color is interpolated. In the case of a classical "traffic light" colormap - a color between yellow and green means "pretty good" and between yellow-red is "pretty bad".
The biggest complication in these scenarios is that it's not always a linear mapping - could be inverse, could be logarithmic, etc. Initially it could support only linear interpolation, but later other types of "library" functions and finally make it completely pluggable.
| Generating rules with an LLM | ||
| ============================== | ||
|
|
||
| Writing a :class:`~airflow.providers.common.dataquality.rules.RuleSet` by hand for every table doesn't scale. |
There was a problem hiding this comment.
There's a middle ground between "by hand" and "via LLM" - deterministic scripts. This also allows users w/o LLM access to use the tools. These can be bundled into the SKILL or placed somewhere in the plugin.
| ============================ | ||
|
|
||
| Quality rules can travel with the :class:`~airflow.sdk.Asset` they describe instead of being | ||
| scattered across every Dag that checks it, and a downstream consumer Dag can refuse to run when |
There was a problem hiding this comment.
I'm not a big fan of describing random bad practices that the code replaces. I suggest
Quality rules can travel with the :class:
~airflow.sdk.Assetthey describeinstead of being scattered across every Dag that checks it, and a downstream consumer Dag can refuse to run when
|
|
||
| Only one Dag should call ``asset_quality()`` for a given asset ``name``/``uri``. Airflow keeps one | ||
| shared record per asset across all Dags, so if more than one Dag attaches (or omits) config for the | ||
| same asset, whichever Dag parsed most recently determines what's stored. |
There was a problem hiding this comment.
Why is this preferred over a loud dag parse error?
| the TaskFlow API. ``ruleset`` may be declared as a decorator argument when it exists at | ||
| Dag-parse time, or returned by the decorated function as a runtime ruleset. ``table`` or | ||
| ``asset`` are declared as decorator arguments exactly like the plain operator. The decorated | ||
| function is optional plumbing on top: return ``None`` to run the check exactly as declared. |
There was a problem hiding this comment.
| function is optional plumbing on top: return ``None`` to run the check exactly as declared. | |
| function is optional: return ``None`` to run the check exactly as declared. |
| If quality checks are one step inside a larger task, use | ||
| :func:`~airflow.providers.common.dataquality.execution.run_quality_checks` directly, then call | ||
| :func:`~airflow.providers.common.dataquality.execution.persist_quality_results` to write the same | ||
| history records that ``DQCheckOperator`` writes: |
There was a problem hiding this comment.
Keeping with existing terminology, wouldn't that be known as a "hook"?




Adds a new
apache-airflow-providers-dqprovider,DbApiHook-based data quality checks.Airflow already has SQL check operators, and many users rely on them for data
quality today. This provider does not replace that path; it adds a small
DQRule/RuleSetlayer for checks that need stable rule identity, persistedhistory, and a connection to Airflow assets. That makes quality results easier
to inspect over time, lets downstream asset consumers gate on recent quality,
and also gives LLM-assisted workflows one schema to generate when proposing
checks from table context. Execution still goes through existing
DbApiHookconnections.
Ships:
DQRuleandRuleSetmodels for named data quality rules.common.sql/DbApiHook, pluscustom_sqlfor database-specific or morecomplex checks.
DQCheckOperatorand the@task.dq_checkTaskFlow decorator.[dq] results_pathfor task, run, andrule-level history.
history.
asset_quality()andrequire_quality(), thatattach provider-owned quality metadata to assets without changing Airflow
core.
LLM-generated rules.
This first version is intentionally small. It focuses on a deterministic rule
shape, SQL execution through
common.sql, persisted results, and lightweightvisibility in the Airflow UI.
Design decisions:
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/local files via
[dq] results_path.The backend stores keyed JSON records for task runs, task instances, and
per-rule history so the UI can read common views without scanning unrelated
runs.
not by changing Airflow core. Static quality configuration is attached to
Asset.extra["airflow.dq"]; runtime summaries are attached to asset eventsunder
extra["airflow.dq.result"]. This lets users try asset quality gatingnow, while leaving room to discuss deeper asset integration later if the
provider gets traction.
DbApiHook/ SQL execution because Airflowalready has strong provider coverage through
common.sql. File andobject-store data checks are left for a later iteration.
later iteration:
files or other object stores and runs quality rules directly against that
data. This PR deliberately starts with the
DbApiHookpath first.Was generative AI tooling used to co-author this PR?
Generated-by: following the guidelines
{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.