Skip to content

Databricks: detect Delta table version changes (#74195). - #74225

Open
goyaladitay11 wants to merge 2 commits into
apache:mainfrom
goyaladitay11:feat/databricks-delta-table-sensor
Open

goyaladitay11 wants to merge 2 commits into
apache:mainfrom
goyaladitay11:feat/databricks-delta-table-sensor

Conversation

@goyaladitay11

Copy link
Copy Markdown

closes: #74195

Description

Adds explicit Delta table version observation to the Databricks provider as Phase 1 of #74195.

Key additions:

  • UnityTableIdentity**: Dataclass representing a static Unity Catalog table identity with validation and conversion to Airflow Asset (to_asset()`).
  • DatabricksDeltaTableVersionSensor: Sensor that queries Delta history (DESCRIBE HISTORY {table_name} LIMIT 1), detects new version commits, records provenance into XCom and Asset event extra` metadata, and supports automatic asset outlet declaration.
  • DatabricksDeltaTableVersionTrigger`: Trigger providing deferrable async polling with bounded backoff.
  • Supports initial enrollment (capturing current version as baseline on first poke), target_version, and allow_recreation handling for tables dropped and recreated.
  • Documented usage, parameters, and reproducibility note under providers/databricks/docs/operators/delta_table_sensor.rst.
  • Registered sensor and trigger in provider.yaml and get_provider_info.py.

@goyaladitay11

goyaladitay11 commented Oct 5, 2026 •

Copy link
Copy Markdown
Author

Hey @zozo123,

Addressed Phase 1 of #74195 with a fully deferrable DatabricksDeltaTableVersionSensor!

Highlights:

Full support for UnityTableIdentity + table_name
Automatic asset outlet events & XCom provenance
Async polling trigger with bounded backoff
Documented the version reproducibility semantics
PR is ready for review: main...goyaladitay11:airflow:feat/databricks-delta-table-sensor. Excited to hear your thoughts and collaborate on Phase 2!

@zozo123

zozo123 commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Thanks @goyaladitay11 — really appreciate you jumping on Phase 1 of #74195. The shape is right: a deferrable DatabricksDeltaTableVersionSensor + trigger on DESCRIBE HISTORY … LIMIT 1, with enrollment / target_version / allow_recreation, docs, and tests is exactly the first slice we wanted for “someone else wrote the table, wake my Dag.”

A few notes before this is mergeable, in the spirit of keeping Unity assets trustworthy for Airflow + Databricks SQL / jobs users:

  1. Please wait on / rebase onto Notify downstream Dags when COPY INTO writes a Unity (Databricks) table #74191. That PR (already approved) introduces UnityTableIdentity and the databricks://…/catalog/schema/table URI conventions with the case-folding fixes potiuk asked for. Databricks: detect Delta table version changes and produce external asset events #74195 explicitly says to reuse that outcome and not introduce a second identity/URI scheme. This branch re-adds UnityTableIdentity in assets/databricks.py — that will collide and risk two slightly different “same table” assets for downstream Dags. Drop the local copy and import from whatever Notify downstream Dags when COPY INTO writes a Unity (Databricks) table #74191 lands.

  2. Be precise about what the outlet event means. Same line we just settled on Databricks: document and validate explicit Unity assets across SQL and job operators #74196: if the sensor declares outlets, the task-success asset event means “a newer Delta version was observed,” with version + table identity in extra — not a vague “table refreshed,” and not a promise that every intermediate commit was delivered. Authors who need a pinned read still have to VERSION AS OF that observed version themselves (your docs already nod at this; please keep that loud).

  3. Label scope as Phase 1 only. Sensor-inside-a-scheduled-task is good. The continuous external producer / coalescing path in Databricks: detect Delta table version changes and produce external asset events #74195 can stay a follow-up so review stays focused.

Happy to re-review once this sits on top of #74191’s identity types. Thanks again for taking this on.

@goyaladitay11
goyaladitay11 force-pushed the feat/databricks-delta-table-sensor branch from 0b77124 to 9bb3d0e Compare October 5, 2026 16:28
@goyaladitay11

Copy link
Copy Markdown
Author

Thanks for the clear feedback @zozo123!

I have addressed all three points:

  1. Rebased onto Notify downstream Dags when COPY INTO writes a Unity (Databricks) table #74191:Rebased this branch on top of pull/74191/head and dropped our local copy of UnityTableIdentity in assets/databricks.py. The sensor now directly imports and uses UnityTableIdentity from Notify downstream Dags when COPY INTO writes a Unity (Databricks) table #74191.
  2. Precise Outlet Event Semantics:
    • Updated the asset outlet extra payload to include observed_version, version, table_identity (host, catalog, schema, table), operation, and timestamp.
    • Updated the documentation to explicitly emphasize that the task-success asset event means "a newer Delta version was observed", that intermediate commits may be coalesced, and that authors needing a pinned read must explicitly specify VERSION AS OF <observed_version> in their queries.
      3.Labeled Scope as Phase 1: Added explicit Phase 1 labeling across the docs, commit, and PR title/description to make clear that this covers the scheduled-task sensor + deferrable trigger, leaving the continuous external asset producer path for Phase 2.

All 46 unit tests for sensor, trigger, and assets are passing. Ready for your re-review whenever you have a chance!

@zozo123 zozo123 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.

Thanks for the follow up. Observation/coalescing semantics and Phase 1 docs look good. Re-reviewed at 9bb3d0e, four things before approval:

  1. Workspace check (P1). The sensor observes via databricks_conn_id but stamps unity_table.host on provenance and the outlet (L174, L249). Please reuse the COPY INTO guard (here): (hook.host or "").lower() != self.unity_table.host.
  2. Deferral drops config (P2). session_configuration, http_headers, client_parameters, hook_params and query tags reach the sync hook but not the trigger (defer, trigger). Please serialize them into the trigger and use them in _get_hook.
  3. Case folding (P2). UnityTableIdentity lowercases, but L152 compares exactly, so Main.Default.Customers fails. Lowercase the parts first, like COPY INTO.
  4. Pinned read docs (P2). The example reads default XCom as a dict; sync returns None, deferred returns an int. Use ti.xcom_pull(task_ids="sensor", key="delta_table_version")["version"] and fix the return_value note on L88.

Also please rebase onto main so the #74191 commits drop out, and add tests for host mismatch and mixed case table_name. Happy to re-review after.

@goyaladitay11
goyaladitay11 force-pushed the feat/databricks-delta-table-sensor branch from 9bb3d0e to 49f321b Compare October 6, 2026 10:03
@goyaladitay11

Copy link
Copy Markdown
Author

Thanks @zozo123! All four points have been addressed, verified locally (all 1,043 unit tests passing), and pushed:

1.Workspace check (P1): Added the COPY INTO workspace host guard (self.hook.host or "").lower() != self.unity_table.host in execute() and poke() before SQL queries run or defer. Added unit tests verifying rejection on host mismatch and case-insensitive host matching.
2. Deferral config (P2):Serialized session_configuration, http_headers, client_parameters, hook_params, and query_tags into DatabricksDeltaTableVersionTrigger and forwarded them when constructing DatabricksSqlHook in _get_hook(). Added tests verifying config serialization and hook instantiation.
3. Case folding (P2):Table parts are now lowercased first in _resolve_table_name() to handle mixed-case inputs like Main.Default.Customers. Added unit tests for mixed-case 1-part, 2-part, and 3-part table names.
4. Pinned read docs (P2):Updated the documentation example to use ti.xcom_pull(task_ids="sensor", key="delta_table_version")["version"] and updated the XCom key note.
5. Rebase: Rebased onto the latest main (cleanly dropping out the #74191 commits now that it has landed) with all tests passing and Ruff formatting clean.

Ready for your re-review!

@zozo123 zozo123 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.

Thanks, re-reviewed at 49f321b. The four earlier points and the rebase look good. One remaining P2 before approval:

Deferral swallows query errors. Sync _get_latest_version raises AirflowException on TABLE_OR_VIEW_NOT_FOUND or permission failures. The trigger except in run only logs a warning and retries until timeout, so a bad table name can burn the full sensor timeout. Please yield an error TriggerEvent for permanent query failures (missing table, access denied) and keep retry only for transient errors, matching the common SQL trigger. Related: restore baseline_version from the trigger event in execute_complete before _push_provenance so deferred enrollment provenance stays accurate after resume.

Happy to approve right after that.

@goyaladitay11
goyaladitay11 force-pushed the feat/databricks-delta-table-sensor branch from 49f321b to 4c4461c Compare October 6, 2026 16:13
@goyaladitay11

Copy link
Copy Markdown
Author

Thanks @zozo123! Both items have been addressed:

Permanent Query Error Handling in Trigger: Added _is_permanent_error() to detect permanent query failures (TABLE_OR_VIEW_NOT_FOUND, PERMISSION_DENIED, ACCESS_DENIED, etc.). The trigger now immediately yields an error TriggerEvent and exits on permanent failures instead of retrying until timeout, preserving retries only for transient errors. Added unit tests for missing table, permission denied, and transient error retry.
Restore Baseline Version on Resume: In execute_complete(), self._baseline_version is now restored from event["baseline_version"] before calling _push_provenance(), ensuring deferred enrollment provenance accurately reflects the baseline version after task resume. Added a unit test verifying this restoration.
All 48 unit tests passing and Ruff formatting clean. Ready for your re-review and approval!

@zozo123 zozo123 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.

Re-reviewed at 4c4461c. Both remaining P2s look good:

  1. Trigger _is_permanent_error now yields an error TriggerEvent and returns on missing-table / permission failures; only transient errors are retried (with tests for both paths).
  2. Success events carry baseline_version, and execute_complete restores it before _push_provenance (covers deferred enrollment).

Earlier workspace host check, deferral config pass-through, case-folding, and pinned-read docs/delta_table_version XCom key still look intact. LGTM for Phase 1.

@goyaladitay11
goyaladitay11 force-pushed the feat/databricks-delta-table-sensor branch from 4c4461c to bc2fb2e Compare October 6, 2026 16:59
@goyaladitay11

Copy link
Copy Markdown
Author

Thank you so much for the thorough and thoughtful reviews throughout this PR, @zozo123! Really enjoyed working on this.
Once this is merged into main, I'm looking forward to collaborating on Phase 2 (the continuous external producer and coalescing path)!

@goyaladitay11

Copy link
Copy Markdown
Author

@potiuk I’ve rebased this PR onto the latest main. @zozo123 has completed the technical review and approved it. The remaining CI workflows are still waiting on maintainer approval — could you please take a look when convenient? Thanks!

Comment on lines +20 to +34
DatabricksDeltaTableVersionSensor (Phase 1)
===========================================

Use the :class:`~airflow.providers.databricks.sensors.databricks_delta_table.DatabricksDeltaTableVersionSensor`
to monitor a Delta table in Databricks and detect new commit versions. When a newer version is observed,
the sensor succeeds, pushes provenance metadata to XCom, and optionally emits an Airflow Asset event.

The sensor can execute synchronously in standard worker slots or in deferrable mode using
:class:`~airflow.providers.databricks.triggers.databricks_delta_table.DatabricksDeltaTableVersionTrigger`.

.. note::

**Scope (Phase 1 only):** This sensor provides task-level Delta table version observation inside a scheduled
DAG task run. A continuous external asset producer with long-running coalescing across runs is planned as a
follow-up Phase 2 feature.

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.

This is use facing doc.
Users don't care about phases. We should just document what we currently have.
If explaining about future work has value then it's OK to mention it but then please explain why it's important to highlight? every feature of airflow can be enhanced further that is how software works :)

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 not very happy with this doc. Feels very AI generated and not designed to be read by humans.
It feels as if the doc explains the code. parameter by parameter.
There are too many bold statements and too many things marked as important.

@eladkal

eladkal commented Oct 8, 2026

Copy link
Copy Markdown
Contributor

cc @moomindani for databricks team review

@goyaladitay11

Copy link
Copy Markdown
Author

Thanks @eladkal — I’ve updated the docs based on your feedback:

  • removed the Phase 1 / Phase 2 wording from the user-facing docs
  • rewrote the page around the user workflow instead of explaining parameters one-by-one
  • removed the heavy bold/important callouts
  • simplified the examples and kept only the behavior users need
  • preserved the observed-version/coalescing and VERSION AS OF semantics in normal user-facing prose
    I also fixed the related docs/static-check issues in the latest commit. Would appreciate another look when you have time.

Aditya Goyal added 2 commits October 8, 2026 20:03
Add DatabricksDeltaTableVersionSensor and DatabricksDeltaTableVersionTrigger to monitor Delta table commit versions (Phase 1: scheduled-task sensor and deferrable trigger):

- Rebase onto main and reuse UnityTableIdentity from airflow.providers.databricks.assets.databricks.
- Validate Databricks connection workspace host against unity_table host case-insensitively.
- Match Unity Catalog table parts case-insensitively, handling mixed-case table_name inputs.
- Forward and serialize session_configuration, http_headers, client_parameters, hook_params, and query_tags into the deferrable trigger.
- Fail fast on permanent query errors (missing table, permission denied) in deferrable trigger while retrying transient errors.
- Restore baseline_version from trigger event in execute_complete before pushing provenance.
- Emit precise task-success asset event extra metadata: observed Delta version and table identity (documenting that intermediate commits may be coalesced, and authors requiring pinned reads must VERSION AS OF).
- Clarify in documentation that pinned reads require VERSION AS OF via ti.xcom_pull(task_ids="...", key="delta_table_version")["version"].
- Support baseline_version, target_version, and allow_recreation handling.
- Support deferrable execution via async triggerer polling.
- Add provider documentation including reproducibility notes and Phase 1 scope callout.
- Register sensor and trigger in provider.yaml and get_provider_info.py.
- Rewrite delta_table_sensor.rst as standard how-to guide, removing Phase 1/2 notes and fixing typo
- Register delta_table_sensor.rst under Databricks SQL how-to-guide in provider.yaml
- Use compat conf and AirflowFailException / AirflowSensorTimeout in sensor
- Narrow query result types in trigger and sensor to resolve MyPy errors
- Defer private baseline version assignment until post-render execution
- Use version-independent outlet_events accessor in test double
@goyaladitay11
goyaladitay11 force-pushed the feat/databricks-delta-table-sensor branch from 7152e22 to a02cf8f Compare October 8, 2026 14:33

@moomindani moomindani 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.

Three findings, inline below. The first two I reproduced with airflow dags test (Airflow main, with only DatabricksSqlHook.run faked to return a version that goes up by one per call). The third I reproduced on a real workspace.

Separately, the docs should mention cost. Each poke runs a query on the SQL warehouse, so a poke_interval shorter than the warehouse's auto-stop keeps it running. This applies in deferrable mode too.


Drafted-by: Claude Code (Opus 5.5); reviewed by @moomindani before posting


def _evaluate_version(self, current_version: int, table_name: str) -> bool:
"""Evaluate if the current version meets sensor criteria."""
if self._baseline_version is None and self.target_version is None:

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.

With mode="reschedule" and no baseline_version, the sensor never succeeds. The enrolled baseline lives only on the operator instance. Every reschedule runs on a fresh instance, so each poke enrolls again and returns False. In my run the table committed 9 times (versions 5 to 13), the sensor logged Initialized baseline version to N on every poke, and it ended in AirflowSensorTimeout. A retry hits the same reset. Poke mode and deferrable mode are fine, because the baseline either stays in the process or is passed to the trigger.

Either keep the enrolled baseline across reschedules, or reject mode="reschedule" without baseline_version/target_version in __init__ so it fails loudly. A test that runs two separate sensor instances against increasing versions would cover it.

return False

if self.target_version is not None:
if current_version >= self.target_version:

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.

target_version is in template_fields, but unlike baseline_version it is never converted with int(). A templated value such as target_version="{{ params.v }}" renders to a string, and the task fails with TypeError: '>=' not supported between instances of 'int' and 'str'. The trigger receives the same string when deferred. Converting it to int in the same place as baseline_version (and testing it with a templated value) fixes both paths.

return self.table_name

def _get_latest_version(self, table_name: str) -> tuple[int, str | None, str | None]:
sql = f"DESCRIBE HISTORY {table_name} LIMIT 1"

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.

DESCRIBE HISTORY ... LIMIT 1 also returns maintenance commits, so the sensor fires when no data changed. On a real Unity Catalog managed table (Predictive Optimization ENABLE (inherited from METASTORE ...), which is the default for newer accounts), I ran two INSERTs and then OPTIMIZE and VACUUM. The history is 1 WRITE, 2 WRITE, 3 OPTIMIZE, 4 VACUUM START, 5 VACUUM END, so LIMIT 1 returns OPTIMIZE or VACUUM END as the "new version". With predictive optimization these commits arrive on their own schedule, so downstream Dags would run on compaction rather than on the external write that #74195 is about. The trigger (triggers/databricks_delta_table.py:131) has the same query.

Consider skipping maintenance operations (for example, reading a few history rows and taking the newest one whose operation is not OPTIMIZE/VACUUM START/VACUUM END), or at least making which operations count configurable and documenting the default.

This branch has not been deployed

No deployments
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.

Databricks: detect Delta table version changes and produce external asset events

4 participants