Repository navigation
Conversation
QAid: zerobus_write_records_qa
namespace: qa
tasks:
- id: write_records
type: io.kestra.plugin.databricks.zerobus.WriteRecords
host: http://host.docker.internal:9090
endpoint: http://host.docker.internal:9090
authentication:
clientId: qa-client
clientSecret: qa-secret
workspaceId: "1234567890"
catalog: main
schema: events
table: user_events
records:
- userId: usr_001
event: page_view
- userId: usr_002
event: purchase |
jymaire
left a comment
There was a problem hiding this comment.
Thanks a lot for this contribution, @Uttam345! The task matches the scope of #214 well: inline records and from (Ion / JSON lines) are both supported, the file is streamed line by line and sent in bounded chunks so memory stays flat, PAT and OAuth client credentials are handled, and the token request follows the Databricks docs. The 29 unit tests pass locally, and the branch also merges cleanly with the open Unity Catalog PR (#288).
I have one point on OAuth token lifetime for long-running ingestions, see the inline comment.
|
@jymaire Thanks, good catch. The OAuth token is now refreshed before a chunk once it is within 60 seconds of expiry (from expires_in, 3600 if the response omits it), and a 401 on the ingest call requests a new token and retries that chunk once. A 401 is an authentication rejection, so the chunk should not have been written and the retry should not duplicate records; I could not verify that against a live workspace. A PAT is never refreshed or retried. Tests with the local HttpServer: a short expires_in refreshes between chunks, a normal one hits the token endpoint once over three chunks, one 401 recovers, two 401s fail with a clear message, and a PAT 401 fails without calling the token endpoint. A good part of the line count in the new commit is formatter re-wrapping; the logic is in TokenHolder, sendChunkWithRetry and getOrRefreshBearerToken. |
jymaire
left a comment
There was a problem hiding this comment.
Kestra Plugin Code Review
Business Requirements — met
The task still covers the Zerobus WriteRecords scope; the OAuth token lifetime concern is now handled.
Previous Review Follow-up — 1 of 1 addressed
Baseline: a1907d5. Reviewed the PR's own delta (f8cbbd1, 1c5110e); the merge of main only brings base changes.
- ✅
WriteRecords.javaOAuth token reused for the whole execution:TokenHolder.needsRefresh()refreshes 60s before expiry (fromexpires_in, default 3600) viagetOrRefreshBearerToken, andsendChunkWithRetryrequests a new token and retries the chunk once on a 401 for OAuth only (PAT guard restored in 1c5110e). Tests cover shortexpires_in, a single token call over three chunks, one 401 recovering, two 401s failing, PAT 401 without token endpoint hit, and a missingexpires_in. Thread resolved.
Kestra Guidelines — 2 findings
- 🟠
WriteRecords.java:217[Guidelines] Stored Ion file read as text lines - 🟢
WriteRecords.java:377[Guidelines] Bare orElseThrow
Security (OWASP Top 10:2025 + KPS) — 0 findings
✅ No issues found
Performance — 0 findings
✅ No issues found
kestraVersion is 1.3.39, so the plugin must also work on 2.0.x: see the 🟠 Ion finding below.
Verdict: RECOMMEND REQUEST CHANGES, comment only (maintainer decision; the token-refresh feedback is fixed, but the 1.3/2.0 Ion read should be addressed; CI has not run yet)
| .orElseThrow(() -> new IllegalArgumentException("from must be provided")); | ||
| reader = new BufferedReader( | ||
| new InputStreamReader( | ||
| runContext.storage().getFile(URI.create(fromUri)), |
There was a problem hiding this comment.
🟠 [Guidelines] Stored Ion file read as text lines
Problem: from is read with BufferedReader.readLine() then JacksonMapper.ofIon().readValue(line), but on Kestra 2.0 files stored by other tasks are binary Ion, so this breaks there while the plugin targets both 1.3 and 2.0.
Fix: Read the stored file with FileSerde.read(InputStream, Consumer) (or Data.from) and keep the chunking in the consumer, instead of reading text lines.
|
|
||
| private void refreshOAuthToken(TokenHolder tokenHolder, RunContext runContext, String host, String workspaceId, | ||
| String catalog, String schema, String table) throws Exception { | ||
| String clientId = runContext.render(getAuthentication().getClientId()).as(String.class).orElseThrow(); |
There was a problem hiding this comment.
🟢 [Guidelines] Bare orElseThrow
Problem: orElseThrow() without a message supplier (here and in resolveBearerToken) would surface as an opaque NoSuchElementException if it ever fires.
Fix: Pass a supplier such as () -> new IllegalArgumentException("clientId must be provided"), or reuse the already-validated values.
Closes #214
Summary
Adds
io.kestra.plugin.databricks.zerobus.WriteRecords, which pushes records into a Unity Catalog Delta table through the Zerobus Ingest REST API. Records come from an inlinerecordslist or from a Kestra internal storage file (from, Ion or JSON lines) and are sent in chunks. The task returnsrecordsCountand emits arecords.countcounter, including when a chunk fails.Design notes
scope=all-apis, the Zerobusresource, andauthorization_detailsfor the catalog, schema and table). A PAT is sent directly as the Bearer token; I could not confirm from the docs that Zerobus accepts one.endpointreplaces the defaulthttps://<workspaceId>.zerobus.<region>.cloud.databricks.com, for other clouds, proxies and tests.Testing
./gradlew test lintPluginDocs: 29WriteRecordsTestcases against a local JDK HttpServer (token request, chunking, errors, no secret in messages). The integration test is skipped unless the DATABRICKS_* variables are set. For each behaviour I broke the matching line and confirmed a test fails.Built from the Databricks docs; not run against a live Zerobus workspace.
Setup: build the shadow jar, mount it in
/app/plugins, and pointhostandendpointat a mock server that answers the token request with anaccess_tokenand the ingest request with an empty JSON object.Written with help from an AI assistant; I reviewed the code and the tests myself.