Skip to content

Add durability to Databricks Lakeflow - #1223

Merged
yaron2 merged 3 commits into
dapr:mainfrom
yaron2:databricks-1
Sep 29, 2026
Merged

yaron2 merged 3 commits into
dapr:mainfrom
yaron2:databricks-1

Conversation

@yaron2

@yaron2 yaron2 commented Sep 18, 2026

Copy link
Copy Markdown
Member

This PR adds an extension that turns Databricks Lakeflow data into durable business actions without building your own retry-safe orchestration layer.

dapr.ext.databricks fixes this with one function call:

  • One line, register_workflow_sink(...), turns any Lakeflow streaming table into a trigger for a durable Dapr Workflow - retries, timers, human approval steps, and multi-day processes, all running independently of your pipeline.
  • Retry- and refresh-safe by default. Every action gets a deterministic, collision-safe identity, so a Lakeflow micro-batch retry - or a full pipeline refresh - is automatically recognized as "already handled," never
    re-triggered.
  • Decouples fast ingestion from slow business logic. A workflow that waits three days for a human doesn't hold up your streaming pipeline for a second - Lakeflow hands off and moves on.

Signed-off-by: yaron2 <schneider.yaron@live.com>
@yaron2
yaron2 requested review from a team as code owners September 18, 2026 22:48
@codecov

codecov Bot commented Sep 18, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 98.85932% with 3 lines in your changes missing coverage. Please review.
✅ Project coverage is 84.30%. Comparing base (d53bd4a) to head (142752b).

Files with missing lines Patch % Lines
dapr/ext/databricks/batch_handler.py 97.59% 2 Missing ⚠️
dapr/ext/databricks/identity.py 97.95% 1 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main    #1223      +/-   ##
==========================================
+ Coverage   83.93%   84.30%   +0.37%     
==========================================
  Files         123      132       +9     
  Lines       10271    10534     +263     
==========================================
+ Hits         8621     8881     +260     
- Misses       1650     1653       +3     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@CasperGN CasperGN 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, this is a useful integration, and the delivery-semantics write-up is clear. The check-then-schedule flow in scheduling.py and the generation epoch for full refreshes both look right to me. tests/ext/databricks passes (71). A few things before this merges:

Blocking

  1. Composite keys can collide, and the second record is silently dropped. id_fields values are joined with _ (identity.py L120), so ('x_y', 'z') and ('x', 'y_z') both produce orders-s-v1-x_y_z. The second record then finds the first record's instance and is logged as already_existed, so its workflow never runs. Please encode the components unambiguously before sanitizing, for example a JSON array or length-prefixed parts, and add a test with two keys that differ only in where the separator falls.

  2. The fallback identity depends on row order. With no key configured, the ID is <batch_id>-<record_index> (identity.py L169). Spark doesn't guarantee the same row order when a micro-batch is retried, so a retry can give record B the index record A had. B is then skipped as already handled, and A is scheduled again under B's old index. Please either require one of id_field/id_fields/instance_id_factory, or make position-based identity an explicit opt-in whose docs say it's only safe for deterministically ordered sources.

  3. The whole batch is buffered in memory despite the docstring. process submits every row to the executor as fast as the iterator yields them (batch_handler.py L122-L135). ThreadPoolExecutor's queue is unbounded, so max_in_flight limits concurrency but not memory. With max_in_flight=2 and a slow sidecar, all 2000 rows of a test batch were already pulled off toLocalIterator() before the first workflow was scheduled. That contradicts "a micro-batch never has to fit entirely in driver memory at once" (L96). A semaphore around submit (released in a done-callback) would bound it to max_in_flight rows.

Non-blocking

  1. id_fields='order_id' (a bare string) is accepted and treated as the fields o, r, d, … (config.py L65-L66). Rejecting str there would catch an easy mistake.
  2. max_records_per_batch fails after scheduling up to the limit (L125). Every retry repeats the same partial work and fails again, so the pipeline stays stuck until the config changes. That's fine as a design choice, but worth a sentence in the docs.
  3. toLocalIterator fallback. On Spark Connect, toLocalIterator is a generator, so an "unsupported" error would surface on the first next(), not at the call inside the try (L165). If the serverless error you observed really is raised at the call, a short comment saying so would help. Otherwise, moving the first next() inside the guard makes the fallback work on both.

Signed-off-by: yaron2 <schneider.yaron@live.com>
@yaron2
yaron2 requested a review from CasperGN September 29, 2026 03:22

@CasperGN CasperGN 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 142752b. All six points from the earlier pass are addressed, with tests:

  • Composite keys are JSON-encoded (identity.py) instead of _-joined, which closes the collision. Regression test added.
  • The position-based fallback identity now requires an explicit allow_batch_position_identity=True. By default, config construction fails.
  • process() bounds read-ahead with a threading.Semaphore(max_in_flight) released in a done-callback, and MemoryBoundedIterationTests covers it.
  • A bare-string id_fields is rejected.
  • The Spark Connect lazy toLocalIterator() failure is covered, because the first row is now pulled inside the guard.
  • The docs now say that max_records_per_batch fails the same way on every retry.

tests/ext/databricks passes (78), and mypy/ruff format are clean. Two minor notes for a follow-up, not this PR: MemoryBoundedIterationTests depends on timing (the slack is generous), and allow_batch_position_identity=True does nothing when a real key strategy is also set.

@yaron2
yaron2 added this pull request to the merge queue Sep 29, 2026
Merged via the queue into dapr:main with commit a8887e3 Sep 29, 2026
22 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants