Skip to content

feat: deliver events to Databricks through Zerobus - #16

Draft
Zaimwa9 wants to merge 4 commits into
feat/postgres-connectionsfrom
feat/databricks-delivery
Draft

Zaimwa9 wants to merge 4 commits into
feat/postgres-connectionsfrom
feat/databricks-delivery

Conversation

@Zaimwa9

@Zaimwa9 Zaimwa9 commented Oct 8, 2026 •

Copy link
Copy Markdown

Delivers events to Databricks through Zerobus Ingest, alongside ClickHouse.

  • DatabricksWarehouse is registered as databricks, built from the connection's host, workspace ID, region, catalog, schema and service principal credentials.
  • It mints its own sql-scoped OAuth token, narrowed to the events table with the zerobuswrite operation, and passes it to the SDK through a headers provider. Customers can use a SQL-only secret instead of an All APIs one. Tokens are cached per credentials and table until 5 minutes before expiry.
  • One stream per batch: ingest_records_offset blocks and raises, then close() waits for acknowledgement. The worst case stays under the 60s insert limit.
  • Rows map to the 11 event columns. Millisecond timestamps become microseconds, falling back to collected_at; non-string values are stored as compact JSON; unusable rows are skipped and logged.
  • SDK and token errors map to the API's failure kinds and sentences. The secret and token never reach errors or logs.

Depends on Flagsmith/flagsmith#8703 for the databricks warehouse type and connection shape.

How did you test this code?

Unit tests cover the adapter fully, with the error mapping built from messages captured live.

Live, against a Databricks trial workspace on default storage:

  • Delivery with a SQL-only secret, a wrong workspace ID, and a schema the service principal has no grants on.
  • Local end to end: SDK-style POST → ingestion server → Redpanda → this service → Databricks, including a failure routed to the retry topic and a successful retry.
  • Flagsmith's exposures and results computation over 524 delivered events matched the generated data exactly.

Not exercised live: a missing table inside a granted schema, and revoking a grant from a working connection.

Follow-up: split batches above the Zerobus payload limit.

@Zaimwa9

Zaimwa9 commented Oct 9, 2026

Copy link
Copy Markdown
Author

@themis-blindfold review

@Zaimwa9 Zaimwa9 changed the title feat: deliver events to Databricks through Zerobus Ingest feat: deliver events to Databricks through Zerobus Oct 9, 2026
Comment on lines +192 to +196
stream = sdk.create_stream(
table_properties=TableProperties(self.table_name),
options=_stream_options(),
headers_provider=_TokenHeaders(self, token),
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟠 Major · ⚡ Quick win

Pass the SDK's required credential parameters.

Observed: create_stream requires client_id, client_secret, and table_properties, even with a custom headers provider; this call supplies only the latter three keywords. Predicted: every nonempty Databricks delivery would raise TypeError before opening a stream. Pass the stored client ID and secret, then cover the real SDK signature.

Suggested change
stream = sdk.create_stream(
table_properties=TableProperties(self.table_name),
options=_stream_options(),
headers_provider=_TokenHeaders(self, token),
)
stream = sdk.create_stream(
self.client_id,
self.client_secret,
TableProperties(self.table_name),
options=_stream_options(),
headers_provider=_TokenHeaders(self, token),
)

headers_provider=_TokenHeaders(self, token),
)
try:
stream.ingest_records_offset(records)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟠 Major · 🏗️ Heavy lift

Split batches before they exceed Zerobus's message limit.

Observed: up to 5,000 Kafka events for one environment reach this single request, while Zerobus caps a message at 10 MB and this adapter converts an oversized payload into DeliveryError. Predicted: a busy tenant's oversized group would be retried intact and then dropped after the retry budget is exhausted. Chunk mapped records by serialized byte size below the SDK limit and ingest each chunk.

@themis-blindfold

Copy link
Copy Markdown

⚖️ Themis review: 🟠 Fix before merge

Databricks delivery cannot open a stream as written, and oversized same-environment batches will be retried unchanged until they are dropped. Fix both before shipping. No CI checks had completed for this revision.

Area Score
🎯 Correctness 1/5
🧪 Test coverage 2/5
📐 Code quality 3/5
🚀 Product impact 4/5

🟠 Majors

  • warehouse_delivery/warehouses/databricks.py:192 — Databricks delivery calls create_stream without its required credential arguments.
  • warehouse_delivery/warehouses/databricks.py:198 — Oversized same-environment groups are retried unchanged and then dropped.
📝 Walkthrough
  • Warehouse registration - registers the Databricks adapter alongside ClickHouse.
  • Authentication - mints a table-scoped OAuth token and supplies it through custom SDK headers.
  • Event delivery - maps accepted event fields to the target table and sends each environment's records through one stream.
  • Tests - cover validation, token failures, record mapping, and SDK error translation, but mock the incompatible stream invocation and do not cover payload chunking.
🧪 How to verify
  • Run uv run pytest tests/unit/warehouses/test_databricks.py after installing the locked dependencies.
  • Exercise create_stream against the real pinned SDK with a custom headers provider and assert the required client ID and secret positions are accepted.
  • Send a same-environment batch whose mapped JSON exceeds the Zerobus message limit and confirm every chunk is acknowledged.
  • Confirm an individual record larger than the provider limit is surfaced as a clear delivery failure without retrying an unchanged oversized batch.
  • Automate: add unit coverage for the real create_stream calling convention and byte-bounded chunking.

Product take: Databricks is a major new warehouse destination, but the current implementation prevents normal delivery and can discard large batches after retries.

🧭 Assumptions & unverified claims
  • The exact pinned SDK wheel could not be run locally because the required Python runtime could not be downloaded; the create_stream signature was checked against Databricks' published Python SDK API.

A stream needs its credentials before it can stream. · reviewed at 26a2cd4

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant