Skip to content

feat(connectors): add generic JDBC source connector - #4272

Open
shbhmrzd wants to merge 37 commits into
apache:masterfrom
shbhmrzd:jdbc_source
Open

shbhmrzd wants to merge 37 commits into
apache:masterfrom
shbhmrzd:jdbc_source

Conversation

@shbhmrzd

Copy link
Copy Markdown

Which issue does this PR address?

Relates to #2500.

This supersedes #3588, which was closed by the stale bot. It includes the review follow-ups from that PR and merges the latest master.

Rationale

Iggy currently has database-specific connectors. A generic JDBC source lets Iggy read from PostgreSQL, MySQL, Oracle, SQL Server, H2, and other JDBC-compliant databases through their existing JDBC drivers.

What changed?

  • Adds a JDBC source connector backed by an embedded JVM and the standard java.sql API.
  • Supports bulk and incremental polling, including persisted offsets through {last_offset} substitution.
  • Maps JDBC values to JSON while preserving precision for decimal and BIGINT values and encoding binary values as base64.
  • Classifies SQL failures by SQLState and clears pending Java exceptions safely.
  • Adds connector configuration, per-database examples, release artifact wiring, and PostgreSQL integration coverage.
  • Updates the shared connector test harness for current master, including persisted topic durability required by source connectors.

How to review

Area Where to look
Source logic and JNI lifecycle core/connectors/sources/jdbc_source/src/lib.rs
Configuration and usage core/connectors/sources/jdbc_source/README.md and config.toml
Integration coverage core/integration/tests/connectors/jdbc/
Shared harness compatibility core/integration/tests/connectors/mod.rs

Verification

  • cargo fmt --all
  • cargo sort --no-format --workspace
  • cargo check -p iggy_connector_jdbc_source
  • cargo clippy -p iggy_connector_jdbc_source --all-features --all-targets -- -D warnings
  • cargo test -p iggy_connector_jdbc_source --all-features (65 passed)
  • cargo clippy -p iggy_connector_sdk --all-targets --all-features -- -D warnings
  • cargo test -p iggy_connector_sdk --all-features (220 passed)
  • cargo clippy -p integration --all-targets --all-features -- -D warnings
  • IGGY_TEST_CLUSTER_NODES=1 cargo test -p integration -- connectors::jdbc:: (7 passed)

The local three-node integration harness did not reach VSR mesh readiness on this machine, before any JDBC test ran. The JDBC suite passes in the harness-supported single-node mode above.

Limitations

  • JNI permits one JVM per process, so JDBC source and sink connectors should run in separate connector-runtime processes.
  • JDBC calls are synchronous on the runtime worker thread.
  • A JVM and database-specific JDBC driver JAR are required at runtime.

AI Usage

Claude Code was used substantially for implementation, tests, documentation, and review hardening. The changes were verified with the commands above and can be explained line by line.

shbhmrzd and others added 24 commits July 20, 2026 20:48
Adds a JDBC source connector that polls a configured SQL query against any
JDBC-compliant database (PostgreSQL, MySQL, Oracle, SQL Server, H2) through an
embedded JVM and produces each row as a JSON message. Supports bulk and
incremental (offset-tracked) modes with persisted offset state.

Also adds the shared connector integration-test harness and a small TCP
listener change so that harness can read the server's written runtime config
when bound to a fixed port. The matching JDBC sink connector follows in a
separate PR.
Aligns the JDBC source crate with the other source connectors: crate
version/type and publish flag, and poll_interval parsing via the workspace
humantime crate rather than a direct humantime-serde pin.

Removes a redundant per-column JNI getObject probe (using wasNull after the
typed getter instead) and caps the per-poll liveness check at five seconds so a
dead connection cannot stall a shared worker. Drops the unused partition_id
field from the example configs and docs.
Derives the incremental cursor from the last ordered row rather than a
Rust-side max, so the cursor always matches the database's own ORDER BY and
degrades to re-reads rather than skips if the ordering is imperfect. Fails
loudly when the tracking column is absent from the result set, and rejects an
empty query, a blank tracking_column, or a blank initial_offset at open().

Routes date/time columns through the null-safe string path so a NULL value no
longer fails the poll, and preserves companion WHERE conditions (such as an
operator-added IS NOT NULL) when stripping the offset predicate at cold start.
…C examples

The incremental examples used timestamp tracking columns (updated_at,
OrderDate) without noting that a non-unique tracking value can skip rows when a
same-value group is split across a batch boundary. Add that caveat to the
affected examples and point to the unique/strictly-increasing requirement.
Incremental-mode validation and row reading disagreed about what a
tracking-column name refers to, so a config that reads rows correctly
could be rejected at open(), and a config that skips rows could be
accepted.

Read-time tracking_column_matches compares against both the raw driver
label and the snake_cased output key, but the ORDER BY check compared a
literal lowercased string. With snake_case_columns on, a snake_cased
tracking_column against a CamelCase ORDER BY column (order_date vs
OrderDate) failed validation despite matching at runtime. The check now
normalizes the ORDER BY key the same way and routes both paths through
tracking_column_matches, and it strips a surrounding identifier quote so
a case-preserving column such as PostgreSQL's ORDER BY "OrderDate"
validates against its unquoted driver label.

Validation also derived the outer ordering with rfind("order by"), which
matched a window function's internal ORDER BY (ROW_NUMBER() OVER (ORDER
BY id) with no outer clause). That orders values within the frame, not
the emitted ResultSet, so setMaxRows still returns an arbitrary subset
and the advancing cursor skips rows, reopening the skip-row bug at
validate time. Ordering detection now scans at parenthesis depth zero,
ignoring ORDER BY inside a subquery, CTE, or OVER (...), and skips
single-quoted string literals so a parenthesis in a literal cannot
perturb the depth. A query whose only ordering sits inside such a
construct is rejected.

Update the README ORDER BY caveat to match and drop a stray line-ending
backslash in the bulk example query.
open() booted the JVM and ran DriverManager.getConnection synchronously
and unbounded, so a single unreachable or slow-DNS database hung the
whole connectors runtime at startup (sources open sequentially, and the
FFI drives open() via block_on), taking every other source, sink, and
the HTTP control API with it. Wrap the JVM init and connection in
block_in_place like poll()/close(), and add login_timeout_ms
(DriverManager.setLoginTimeout) so a stuck connect fails instead of
blocking forever.

BIGINT was emitted as a JSON number, silently rounding values above 2^53
in consumers that parse JSON numbers as f64. Emit it as a string, the
same as NUMERIC/DECIMAL already do.

The incremental cursor comment claimed it "degrades to re-reads rather
than skips", which is false when a run of rows sharing one tracking
value is split across a setMaxRows batch boundary: the remainder is
skipped. Correct the comment and warn at validate time that the tracking
column must be unique or batch_size must exceed the largest tie group.

Hoist the jni dependency to the workspace so the pin is shared with the
upcoming JDBC sink rather than duplicated.

Document that SecretString redaction is not end-to-end: the runtime's
generic configs/plugin control endpoint and trace-level config logging
emit the raw plugin_config for every connector.
# Conflicts:
#	Cargo.lock
#	core/integration/tests/connectors/mod.rs
A proactive claims-vs-code sweep found two README statements the code
does not back. `table_name` is a reserved field hardcoded to null, never
derived from the query, yet the metadata was advertised as including the
table name; document that it is always null and that operation_type is
always SELECT. "Prevents duplicate reads" contradicted the connector's
own at-least-once guarantee; reword to "avoids re-reading rows below the
tracked offset (at-least-once, not exactly-once)".
Harden outer_order_by() to skip SQL line (--) and block (/* */) comments,
so a commented-out ORDER BY can no longer satisfy the incremental
ordering requirement and let an unordered query pass validation and skip
rows at runtime.

Complete JNI exception clearing in throwable_string_method(): a Java
exception raised by getMessage()/getSQLState() or the string conversion
is now cleared before returning, so the next JNI call on the thread is
not aborted for running with an exception pending.

Drop the Serialize derive (and the serialize_secret helpers) from
JdbcSourceConfig: the runtime only deserializes it, and serializing would
risk writing the jdbc_url/password SecretStrings in plaintext. This also
removes the now-unused iggy_common dependency. Redacted logging still goes
through the manual Debug impl.

Read column names with ResultSetMetaData.getColumnLabel() instead of
getColumnName() so a SELECT expr AS alias yields the alias the caller
configured rather than the base-table column (or an empty string).

Remove the unused connectors_api_address() test helper that failed
integration clippy, and make the JDBC integration tests panic on a
Postgres/testcontainers setup failure instead of returning early and
reporting a false pass.
Correct "unparseable" -> "unparsable" in poll_interval doc comments and
code comments (flagged by the typos check), and reflow jdbc_oracle.toml
to taplo's canonical format (a stray blank line).
The JDBC Postgres tests downloaded the driver JAR at runtime and cached
it without validating it. When the download returned a truncated body or
an error page (seen in CI), the file was still cached and put on the JVM
classpath; the JVM booted fine but Class.forName("org.postgresql.Driver")
then failed, so the connector never initialized and every test expecting
messages saw zero, while the bad jar was reused by every later test.

Download from Maven Central, verify the bytes are a real JAR (ZIP magic
plus a minimum size) before persisting, and write via a temp file and
atomic rename so a partial or corrupt download can never be cached. A
previously cached invalid jar now self-heals, and a genuinely bad
download fails setup with a clear message instead of a later opaque
ClassNotFoundException.
# Conflicts:
#	Cargo.toml
#	core/server/src/tcp/tcp_listener.rs
The incremental cursor advanced as soon as rows were fetched, so a batch
whose send later failed was skipped permanently: the next poll queried
past it, and its further-advanced offset then overwrote the checkpoint,
removing any chance of replay. The connector could not fix this before,
because the runtime handed a source its state only at open() with no
per-batch delivery signal.

Upstream now supplies that signal. Stage the advanced cursor in poll()
and resolve it in on_batch_result: Ack (batch sent and checkpoint
persisted) commits it, Nack discards it so the same range is re-polled.
That makes at-least-once hold for an ordinary in-process send failure,
not just across restarts, and follows the staging pattern documented for
sources and implemented by random_source.

An empty poll now checkpoints nothing rather than rewriting an unchanged
cursor, which on the HTTP state backend costs a round-trip per interval.

Also repair the shared connectors test harness for the create_topic and
send_messages signature changes that came with master, and rename the
JDBC integration test file to the sibling convention.
Only a status code crosses the FFI boundary when a plugin is opened, and
both containers discarded the error behind `result.is_ok()`. The runtime
could then report no more than "plugin initialization failed", so an
operator saw a connector silently skipped at startup with nothing to say
why, and the plugin's own diagnosis was thrown away.

A JDBC source whose driver JAR could not be read cost a full CI
investigation for exactly this reason: the connector had already
produced "Failed to load driver class 'org.postgresql.Driver'" and it
never reached a log.

Log the cause before collapsing it to a code, for sinks and sources
alike. The signature is unchanged, so pre-built plugins keep working.
The driver-JAR integrity check accepted any file of at least 500 KB
starting with the ZIP local-file-header magic. A download truncated past
that size still satisfies both conditions while its central directory is
gone, so it is not readable as an archive at all. The check's own comment
claimed such a download would fail it.

That file was then cached and reused by every later JDBC test in the job,
and the JVM could not load the driver class from it. `Class.forName`
threw, the connector's `open()` failed, and the runtime skipped the
source, so tests failed having received no messages rather than pointing
at the JAR. One bad download made this deterministic across all four
nextest retries and took out every JDBC test sharing the runner.

Verify instead that the bytes open as an archive containing
`org/postgresql/Driver.class`, which is the property the JVM needs; a
cached file that fails is already deleted and re-downloaded.

Also stop comparing the recovery assertion against the raw message list.
Bulk mode re-runs its query every poll interval, so it re-delivers the
same rows indefinitely and two batches arriving in one client poll turned
a healthy run into an id mismatch. Match on the distinct ids seen, and
report the received count so "nothing was delivered" is distinguishable
from "delivered without the expected data.id".
The shared connectors test client polled in batches of 10, and the JDBC
collect-until-N loop gave up after 30 attempts. Those multiply into a
ceiling: the loop could never observe more than 300 messages however many
the source delivered. large_result_set_streams_without_crashing waited for
150 of that 300 and locally spent ~16 attempts draining 10 at a time, so a
slower runner exhausted the budget before the count was reached and the
test failed on every retry rather than flaking. The five tests that pass
are all bulk mode, which re-runs its query each interval and so retries
for free; the two that failed were the only ones whose expectations could
not absorb a slow runner.

Let the caller choose the batch count, request more per poll than any test
waits for, and bound the wait with a wall-clock deadline so runner speed
costs time instead of correctness.

Assert the incremental cursor on the distinct ids seen rather than on a
fixed-length prefix of the received messages. The claim under test is
unchanged, that the offset advanced past 3, but a re-delivered batch no
longer fails it, which matters because the connector's delivery contract
is at-least-once. Both panics now report the ids seen and the raw message
count, so a future failure distinguishes a stalled source from a replay.
# Conflicts:
#	.config/nextest.toml
#	.github/workflows/_build_rust_artifacts.yml
#	Cargo.lock
#	Cargo.toml
@github-actions

Copy link
Copy Markdown

Thanks for the PR. It is labeled S-waiting-on-review and queued for review.

Slash commands (own line, regular comment) move it around the queue:

  • /ready - back to S-waiting-on-review after addressing feedback
  • /author - flip to S-waiting-on-author while you finish changes
  • /request-review @user-or-team - request a reviewer
  • /pin - exempt the PR from the stale bot, /unpin to undo

See CONTRIBUTING.md for details.

@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Sep 23, 2026
@github-actions
github-actions Bot requested a review from hubcio September 23, 2026 13:55
@codecov

codecov Bot commented Sep 23, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 70.00000% with 3 lines in your changes missing coverage. Please review.
✅ Project coverage is 87.70%. Comparing base (15a49d0) to head (15be65e).
⚠️ Report is 16 commits behind head on master.

Files with missing lines Patch % Lines
core/connectors/sdk/src/sink.rs 40.00% 3 Missing ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master    #4272      +/-   ##
============================================
+ Coverage     87.56%   87.70%   +0.14%     
- Complexity     1575     1576       +1     
============================================
  Files          1284     1285       +1     
  Lines        225457   227541    +2084     
  Branches     188821   190905    +2084     
============================================
+ Hits         197413   199570    +2157     
+ Misses        23309    23154     -155     
- Partials       4735     4817      +82     
Components Coverage Δ
Rust Core 88.80% <70.00%> (+0.04%) ⬆️
Java SDK 68.70% <ø> (+0.01%) ⬆️
C# SDK 77.41% <ø> (-0.02%) ⬇️
Python SDK 90.97% <ø> (ø)
PHP SDK 85.67% <ø> (ø)
Node SDK 96.42% <ø> (+1.67%) ⬆️
Go SDK 70.15% <ø> (+0.10%) ⬆️
Files with missing lines Coverage Δ
core/connectors/sdk/src/source.rs 90.51% <100.00%> (+0.64%) ⬆️
core/connectors/sources/jdbc_source/src/lib.rs 91.18% <ø> (ø)
core/connectors/sdk/src/sink.rs 80.07% <40.00%> (-0.48%) ⬇️

... and 43 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@hubcio

hubcio commented Sep 26, 2026

Copy link
Copy Markdown
Contributor

you dont need to constantly rebase. maintainers will do that :)

@hubcio

hubcio commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

was I not clear enough? 🥲 please, dont rebase, your PR is merge-able without conflicts, it's good for review now, expect review this week. if it won't rebase cleanly we'll ping you.

@shbhmrzd

shbhmrzd commented Sep 28, 2026 •

Copy link
Copy Markdown
Author

was I not clear enough? 🥲 please, dont rebase, your PR is merge-able without conflicts, it's good for review now, expect review this week. if it won't rebase cleanly we'll ping you.

My bad, there was a conflict earlier and thus had to merge and resolve conflicts.

@hubcio

hubcio commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

/skill team-review-slim

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Summary: The JDBC source connector's staged-cursor protocol, JNI error clearing and SQL validation hold up under tracing; fifteen findings remain, none critical. The strongest are a zero poll_interval that busy-loops the database, an incremental query without a cursor predicate that re-reads the same rows forever, a tie group that can skip rows, and two credential leaks around the logs - a URL sanitizer that stops at the first @​, and the new SDK open-error log that can print the full jdbc_url; claims this checkout cannot settle, such as jni-rs method lookups and non-MySQL backslash escaping, were dropped.

Counts: critical 0, warning 8, nit 3, simplification 4

Findings without an anchor on a changed line:

  • core/connectors/sources/jdbc_source/src/lib.rs:1536 warning: poll_interval = "0s" parses, so the sleep never runs and poll() issues back-to-back JDBC queries against the database. Reject a zero interval in validate_config, or clamp it to a 1ms floor.
  • core/connectors/sources/jdbc_source/src/lib.rs:968 warning: read_rows takes the last row's tracking value as the cursor, so a tie group split by setMaxRows loses its remaining rows on the next strict > poll. Probe one row past batch_size in incremental mode and fail closed when the extra row repeats the last in-batch value.
  • core/connectors/sources/jdbc_source/src/lib.rs:1450 warning: Incremental mode accepts a query without {last_offset}, so every poll re-reads the same first batch_size rows while the cursor never enters the query. Require the placeholder in the query at open(), or error the poll when the built query holds no cursor comparison.
  • core/connectors/sources/jdbc_source/src/lib.rs:1027 warning: read_single_row takes the cursor from the first result column whose label matches tracking_column, while query_orders_by_tracking_column strips the ORDER BY qualifier, so a join that returns two updated_at labels can set the cursor from the column the query does not order by. Reject a result set with more than one matching label, or take the cursor from the ordered column.
  • core/connectors/sources/jdbc_source/src/lib.rs:143 warning: sanitize_jdbc_url stops at the first @​, so a jdbc_url whose password contains one, such as user:p@​ss@​host, still logs ss@​host at lines 444 and 1503. Mask the credential up to the last @​ before the host.
  • core/connectors/sources/jdbc_source/src/lib.rs:1006 warning: read_single_row rebuilds each column's output key and repeats the collision warning for every row, so one collapsed key pair logs once per row per poll. Compute the key and warn once per column, in read_column_metadata.
  • core/connectors/sources/jdbc_source/src/lib.rs:1702 warning: get_or_create_jvm returns the live JVM without comparing driver_jar_path or jvm_options, so a second JDBC source silently uses the first connector's classpath. Fail at open() when the requested classpath differs, because validate_config already accepted the unused JAR path.
  • core/connectors/sources/jdbc_source/src/lib.rs:367 simplification: state is an Arc<​Mutex<​State>> that nothing clones, because the source is its only owner. Simpler: use a plain Mutex<​State>, which removes one heap allocation and one level of indirection.
  • core/connectors/sources/jdbc_source/src/lib.rs:334 simplification: State.last_poll_time is written into every checkpoint and never read. Simpler: delete the field, since the connector is new and no older checkpoint depends on its position.
  • core/connectors/sources/jdbc_source/src/lib.rs:2044 simplification: java::sql::Types defines CHAR, VARCHAR, LONGVARCHAR and NULL, and nothing reads them, so two #[allow(dead_code)]​ attributes hide the gap. Simpler: delete the four constants and both attributes.

This review was generated by Claude Code 2.1.284 on deepseek-flash[1m]. Review the output before you act on it.

Comment thread core/connectors/sdk/src/source.rs Outdated
match result {
Ok(()) => 0,
Err(error) => {
error!("Failed to open source connector with ID: {id}. {error}");

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

warning: The new log line writes the connector's open error verbatim, and a driver rejection such as JDBC's No suitable driver found for <​url> echoes the full jdbc_url, credentials included. Sanitize the URL out of that error before the connector returns it. Also at core/connectors/sdk/src/sink.rs:101.

```sql
-- Configuration
tracking_column = "id"
query = "SELECT * FROM users WHERE id > {last_offset} ORDER BY id"

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

nit: The example config has no initial_offset, and strip_offset_predicate removes only the literal WHERE {tracking_column} > {last_offset} form, so this query fails open() on the first poll. Add initial_offset = "0" to the example.

password = ""

# Simple query for testing
query = "SELECT * FROM users WHERE id > {last_offset} ORDER BY id"

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

nit: The in-memory H2 database starts empty, and nothing creates the users table the query selects, so this example errors on every poll. Add an INIT clause that creates and seeds the table, as test_jdbc_h2.toml does.

uuid = { workspace = true, features = ["v4"] }

[dev-dependencies]
toml = { workspace = true }

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

nit: The crate omits [lints]​ workspace = true, so the workspace warnings = "deny" gate does not cover it. Add the section, as the five sibling source crates do.

impl ConnectorsIggyClient {
/// Send messages to the configured stream/topic (used by sink connector tests).
#[allow(dead_code)]
async fn send_messages(

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

simplification: ConnectorsIggyClient::send_messages is never called, and its doc comment claims a sink test uses it, so #[allow(dead_code)]​ hides 15 lines of dead test scaffolding. Simpler: delete the method.

@github-actions github-actions Bot added S-waiting-on-author PR is waiting on author response and removed S-waiting-on-review PR is waiting on a reviewer labels Sep 28, 2026

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

S-waiting-on-author PR is waiting on author response

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants