Repository navigation
Conversation
|
@themis-blindfold review |
| stream = sdk.create_stream( | ||
| table_properties=TableProperties(self.table_name), | ||
| options=_stream_options(), | ||
| headers_provider=_TokenHeaders(self, token), | ||
| ) |
There was a problem hiding this comment.
🟠 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.
| 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) |
There was a problem hiding this comment.
🟠 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 review: 🟠 Fix before mergeDatabricks 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.
🟠 Majors
📝 Walkthrough
🧪 How to verify
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
A stream needs its credentials before it can stream. · reviewed at 26a2cd4 |
Delivers events to Databricks through Zerobus Ingest, alongside ClickHouse.
DatabricksWarehouseis registered asdatabricks, built from the connection's host, workspace ID, region, catalog, schema and service principal credentials.sql-scoped OAuth token, narrowed to the events table with thezerobuswriteoperation, 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.ingest_records_offsetblocks and raises, thenclose()waits for acknowledgement. The worst case stays under the 60s insert limit.collected_at; non-string values are stored as compact JSON; unusable rows are skipped and logged.Depends on Flagsmith/flagsmith#8703 for the
databrickswarehouse 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:
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.