Skip to content

feat(ingest): add zerobus WriteRecords task - #289

Open
Uttam345 wants to merge 5 commits into
kestra-io:mainfrom
Uttam345:feat/zerobus-write-records
Open

Uttam345 wants to merge 5 commits into
kestra-io:mainfrom
Uttam345:feat/zerobus-write-records

Conversation

@Uttam345

@Uttam345 Uttam345 commented Oct 2, 2026 •

Copy link
Copy Markdown

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 inline records list or from a Kestra internal storage file (from, Ion or JSON lines) and are sent in chunks. The task returns recordsCount and emits a records.count counter, including when a chunk fails.

Design notes

  • OAuth client credentials use the token request from the Databricks REST docs (scope=all-apis, the Zerobus resource, and authorization_details for 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.
  • An optional endpoint replaces the default https://<workspaceId>.zerobus.<region>.cloud.databricks.com, for other clouds, proxies and tests.
  • The chunk limits (500 records, 4 MB) are my assumption: the docs give 10 MB per record and no per-request limit.
  • Delivery is at-least-once, so a task retry can write records twice. When a chunk fails the task throws and reports how many records were already accepted.
  • An OAuth token is refreshed shortly before it expires, and a 401 on the ingest call triggers one refresh and one retry of that chunk. A PAT is never refreshed or retried.
  • The icon is a copy of the Databricks one.

Testing

./gradlew test lintPluginDocs: 29 WriteRecordsTest cases 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 point host and endpoint at a mock server that answers the token request with an access_token and the ingest request with an empty JSON object.

Screenshot 2026-10-03 at 02 59 42

Written with help from an AI assistant; I reviewed the code and the tests myself.

@Uttam345

Uttam345 commented Oct 2, 2026 •

Copy link
Copy Markdown
Author

QA

id: 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

@fdelbrayelle
fdelbrayelle requested review from a team and Malaydewangan09 October 4, 2026 11:37
@fdelbrayelle fdelbrayelle added area/plugin Plugin-related issue or feature request kind/external Pull requests raised by community contributors labels Oct 4, 2026

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

Comment thread src/main/java/io/kestra/plugin/databricks/zerobus/WriteRecords.java Outdated
@Uttam345

Uttam345 commented Oct 6, 2026

Copy link
Copy Markdown
Author

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

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.java OAuth token reused for the whole execution: TokenHolder.needsRefresh() refreshes 60s before expiry (from expires_in, default 3600) via getOrRefreshBearerToken, and sendChunkWithRetry requests a new token and retries the chunk once on a 401 for OAuth only (PAT guard restored in 1c5110e). Tests cover short expires_in, a single token call over three chunks, one 401 recovering, two 401s failing, PAT 401 without token endpoint hit, and a missing expires_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)),

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.

🟠 [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();

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.

🟢 [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.

@fdelbrayelle
fdelbrayelle removed the request for review from Malaydewangan09 October 9, 2026 10:01
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area/plugin Plugin-related issue or feature request kind/external Pull requests raised by community contributors

Projects

Status: To review

Development

Successfully merging this pull request may close these issues.

feat(ingest): add Zerobus Ingest subplugin for push-based ingestion into Unity Catalog Delta tables

3 participants