Skip to content

feat(connectors): add Apache Fluss source connector - #3799

Open
seokjin0414 wants to merge 6 commits into
apache:masterfrom
seokjin0414:3688-fluss-source
Open

seokjin0414 wants to merge 6 commits into
apache:masterfrom
seokjin0414:3688-fluss-source

Conversation

@seokjin0414

@seokjin0414 seokjin0414 commented Aug 2, 2026 •

Copy link
Copy Markdown
Contributor

What

Adds a source connector that reads an Apache Fluss
log table and publishes each row into an Apache Iggy stream as JSON. This is
phase 1 of #3688: log tables and Schema::Json, on the fluss-rs 1.0 client.

How it works

The connector opens one table snapshot and takes the column list, the bucket
count and the log scanner from it, so a missing table, a primary-key table, a
partitioned table or a column type with no JSON form all fail at startup rather
than once per batch. The scanner aligns every row to the schema it was created
with, so decoding against that same snapshot keeps columns in place when the
table's schema changes under a running connector. The scanner is created once
in open() and reused, which works because it owns its client state rather
than borrowing the connection.

Offsets are tracked per bucket and follow the batch acknowledgment contract
(#3855): poll stages the candidate offsets with the batch and
on_batch_result commits them on Ack. The scanner moves past records as soon
as it returns them, so a Nack rewinds it to the committed offsets and those
rows are read again instead of skipped. Fetches still buffered from the old
position are dropped by the scanner's own expected-offset check. The rewind
retries, because it can need a metadata refresh and an error from
on_batch_result stops the source.

A row the JSON mapping cannot hold, such as a date beyond what chrono
represents, fails the same way on every read. Rewinding for it would stall the
source for good, and since the SDK only logs a poll error, the source would
still look healthy. Such a row is dropped and logged at error level with its
bucket and offset, and the offsets move past it. The runtime and the HTTP sink
handle messages they cannot decode the same way. A plugin cannot bump the
runtime's error metric, so close() logs how many rows were skipped.

An empty poll carries no state, with one exception: the first poll after
open() persists the resolved start offsets. Without it, a restart before the
first row would resolve latest again and skip whatever was written in
between. On restart, buckets already present in the state keep their offset
and only new ones fall back to starting_offset, which keeps a widened bucket
count from rewinding buckets that were already consumed. Only the table's
current buckets are subscribed, so a saved offset for a bucket the table no
longer has is ignored with a warning.

Each message id is derived from its bucket and offset, so a consumer can spot a
record replayed after an at-least-once redelivery (the server does not dedupe
on it). Column projection is pushed down to the server when columns is set.
The names are resolved to positions once and both the row decoder and the
scanner use them, so a projection listed in a different order than the table
still maps each value to its own column. With include_metadata on, the
bucket, offset and timestamp are added under the _fluss_ prefix, so a table
column under that prefix is rejected at startup instead of being overwritten.

Temporal values are formatted with every fractional digit the column holds, the
way the PostgreSQL source formats them: TIMESTAMP without a timezone,
TIMESTAMP_LTZ as RFC 3339 in UTC. A number of milliseconds would drop the
microseconds of the default TIMESTAMP(6).

Settings that would otherwise fail late or silently are rejected before
connecting:

  • empty bootstrap_servers, database or table
  • an empty or repeated columns list
  • a zero batch_size, which the client accepts and which makes every poll come
    back empty
  • a negative starting_offset
  • only one of the two SASL credentials
  • a zero poll_interval together with a zero poll_timeout, which would poll
    an idle table in a tight loop

The protoc question from #3688

Resolved upstream: apache/fluss#3874 checked the generated protobuf code into
fluss-rs, and 1.0.0 is the first release with it. This PR is now a plain
workspace member with no system dependency and no CI change, and the earlier
ci: commit is gone.

One dependency note: fluss-rs 1.0 uses arrow 59 while the workspace stays on
58 for iceberg and deltalake, so the lockfile carries a second arrow line until
those catch up.

Testing

Unit tests cover config validation, the SASL client config, the state round
trip, offset resolution (including buckets the table no longer has), the
ACK/NACK staging (including the start offsets riding an empty batch and the
offsets past skipped rows), the message metadata, the projection lookup and the
row-to-JSON mapping (temporal formatting down to nanoseconds and before the
epoch, decimals, nulls, base64 for binary, non-finite floats). I broke each
staging rule on purpose and checked that a test fails.

The integration tests run against a real Fluss 1.0.0 cluster:

  • Rows appended to a log table come out of the Apache Iggy topic with their
    offsets and bucket.
  • With the Apache Iggy server killed, a batch is rejected after the scanner has
    handed its rows over. Once the server is back every row still arrives, which
    only holds if the source rewound the scanner.
  • Starting from latest over a table that already holds rows, the runtime is
    restarted before any new row arrives. The rows written while it was down are
    delivered and the older ones are not.
  • A row holding every supported column type, plus a row of nulls, comes out as
    the expected JSON. These rows are decoded from Arrow by the server, the path
    production reads, while the unit tests build rows in memory.
  • A row holding a date chrono cannot represent sits between two good rows.
    The good rows arrive and the bad one is skipped.
  • A projection names two columns in reverse table order and leaves a third out.
    Each value arrives under its own name and the third column is absent.
  • Read without a projection, that same table has a column under _fluss_, so
    with include_metadata on the source fails to start. The projection test
    reads the table fine, so the failure comes from that column.

Each test for a fix fails when the fix is taken out. Without the rewind no row
arrives after the restart, without persisting the start offsets the state file
never appears, and with the old millisecond timestamps the TIMESTAMP(6)
column comes back as a number. Failing the batch on the bad row delivers
nothing, decoding in table order instead of projection order delivers no row,
and without the prefix check the source starts.

Locally, prek run for both the pre-commit and pre-push stages (fmt, sort,
clippy with --all-features, taplo, typos, license headers and the rest),
rustdoc with -D warnings, the unit tests and the seven integration tests all
pass.

The cluster runs as a single container. The image ships local-cluster.sh,
which starts the same embedded ZooKeeper, coordinator and tablet server, but it
rewrites the tablet server's bind port to 0 so it cannot collide with the
coordinator, and a random port inside the container cannot be published to the
host. Giving the tablet server an explicit second port keeps both reachable and
lets both advertise localhost, which then resolves the same way inside the
container and from the test process.

Notes on scope

These are limits of the connector, each rejected at startup with the reason
rather than silently ignored:

  • Primary-key tables. fluss-rs 1.0 can read their changelog, but mapping its
    insert, update and delete records onto messages deserves its own PR.
  • payload_format = "arrow_ipc". It needs the batch scanner, which has its own
    offset-tracking path, so it is left for a follow-up.
  • Partitioned tables.

The README also lists behavior a user should know about: rows that cannot be
converted are skipped, a crash right after starting from latest can miss
rows, and an offset outside what a bucket still holds is not reset.

Relates to #3688

@github-actions

github-actions Bot commented Aug 2, 2026

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

See CONTRIBUTING.md for details.

@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Aug 2, 2026
@codecov

codecov Bot commented Aug 2, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 92.01878% with 85 lines in your changes missing coverage. Please review.
✅ Project coverage is 87.96%. Comparing base (b8bf85f) to head (65eab7d).
⚠️ Report is 59 commits behind head on master.

Files with missing lines Patch % Lines
core/connectors/sources/fluss_source/src/lib.rs 93.42% 42 Missing and 11 partials ⚠️
...ore/connectors/sources/fluss_source/src/mapping.rs 87.64% 6 Missing and 26 partials ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master    #3799      +/-   ##
============================================
+ Coverage     87.57%   87.96%   +0.39%     
- Complexity     1575     1579       +4     
============================================
  Files          1284     1292       +8     
  Lines        225457   231210    +5753     
  Branches     188821   194569    +5748     
============================================
+ Hits         197438   203385    +5947     
+ Misses        23289    23033     -256     
- Partials       4730     4792      +62     
Components Coverage Δ
Rust Core 89.06% <92.01%> (+0.29%) ⬆️
Java SDK 68.74% <ø> (+0.05%) ⬆️
C# SDK 77.71% <ø> (+0.28%) ⬆️
Python SDK 91.24% <ø> (+0.26%) ⬆️
PHP SDK 85.67% <ø> (ø)
Node SDK 96.59% <ø> (+1.84%) ⬆️
Go SDK 70.26% <ø> (+0.21%) ⬆️
Files with missing lines Coverage Δ
...ore/connectors/sources/fluss_source/src/mapping.rs 87.64% <87.64%> (ø)
core/connectors/sources/fluss_source/src/lib.rs 93.42% <93.42%> (ø)

... and 224 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.

@github-actions

Copy link
Copy Markdown

This pull request has been automatically marked as stale because it has not had recent activity. It will be closed in 7 days if no further activity occurs.

If you need a review, please ensure CI is green and the PR is rebased on the latest master. Don't hesitate to ping the maintainers - either @core on Discord or by mentioning them directly here on the PR.

Thank you for your contribution!

@github-actions github-actions Bot added the S-stale Inactive issue or pull request label Aug 13, 2026
@github-actions github-actions Bot removed the S-stale Inactive issue or pull request label Aug 14, 2026
@slbotbm

slbotbm commented Aug 16, 2026

Copy link
Copy Markdown
Contributor

Would be good to convert this to a draft until upstream changes are merged. This will get autoclosed otherwise

@seokjin0414
seokjin0414 marked this pull request as draft August 20, 2026 15:36
@github-actions github-actions Bot removed the S-waiting-on-review PR is waiting on a reviewer label Aug 20, 2026
@seokjin0414

Copy link
Copy Markdown
Contributor Author

Rebased onto latest master and adapted to the source batch acknowledgments (#3855): offsets are now staged in poll and committed on Ack; a Nack rewinds the scanner to the last acknowledged offsets so a rejected batch is read again instead of skipped. Unit + Docker e2e pass locally, CI is green.

Upstream apache/fluss#3874 landed, so the next fluss-rs release drops the protoc requirement. Keeping this as a draft until then; once it ships I'll drop the two protoc commits and mark this ready for review.

@seokjin0414
seokjin0414 marked this pull request as ready for review September 24, 2026 14:15
@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Sep 24, 2026
@seokjin0414

Copy link
Copy Markdown
Contributor Author

fluss-rs 1.0.0 is out with the generated protobuf code from apache/fluss#3874, so the protoc commits are gone and this is now a plain workspace member with no CI change.

Rebased onto master and picked up the newer source rules on the way (no state on empty polls, a bounded retry on the rewind). Two fixes came out of that pass: the start offsets resolved for latest are now persisted before the first row arrives, and a batch that fails to build rewinds the scanner instead of skipping its rows. Both have a new integration test against Fluss 1.0.0, and CI is green.

@hubcio this is the one the upstream change was for, ready for review when you have time.

@seokjin0414
seokjin0414 force-pushed the 3688-fluss-source branch 2 times, most recently from fa9c181 to 8847622 Compare September 25, 2026 03:22
Comment thread core/connectors/sources/fluss_source/Cargo.toml
Comment thread core/integration/tests/connectors/fluss/fluss_source.rs Outdated
Comment thread core/integration/tests/connectors/fluss/fluss_source.rs
Comment thread core/connectors/sources/fluss_source/src/lib.rs Outdated
Comment thread core/connectors/sources/fluss_source/src/lib.rs Outdated
Comment thread core/connectors/sources/fluss_source/src/lib.rs Outdated
Comment thread core/connectors/sources/fluss_source/src/lib.rs
Comment thread core/connectors/sources/fluss_source/src/lib.rs Outdated
Comment thread core/connectors/sources/fluss_source/src/lib.rs
Comment thread core/connectors/sources/fluss_source/src/lib.rs Outdated
@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 26, 2026
@seokjin0414

seokjin0414 commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor Author

Addressed the review in 828f4d4, rebased onto master, which now has the Floci switch from #4290. Beyond the comments, a zero poll_interval together with a zero poll_timeout is now rejected too, since that would poll an idle table in a tight loop.

/ready

@github-actions github-actions Bot removed the S-waiting-on-author PR is waiting on author response label Sep 27, 2026
@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Sep 27, 2026
Apache Fluss keeps streams as schema-aware columnar log tables, so feeding one
into Apache Iggy so far meant writing a bespoke client. This connector reads a
Fluss log table through fluss-rs 1.0 and publishes each row as JSON, tracking
offsets per bucket through the runtime state API so a restart resumes where the
previous run stopped. Buckets already present in the restored state keep their
offset, which keeps a widened bucket count from rewinding buckets that were
already consumed.

Offsets are staged with each batch and committed only when the runtime
acknowledges it. The scanner moves past records as soon as it returns them, so
a rejected batch, or one that cannot be built, rewinds it to the committed
offsets and those rows are read again rather than skipped. The start offsets
ride the first poll even when it has no rows. Otherwise a restart before the
first row would resolve `latest` again and skip whatever was written in
between.

Primary-key tables stream a changelog whose change types need their own mapping
onto messages, and arrow_ipc payloads need the batch scanner with its own
offset-tracking path, so both are rejected at startup instead of being silently
downgraded. Column projection is pushed down to the server when `columns` is
set.

Temporal values are formatted with every fractional digit the column holds, the
way the PostgreSQL source formats them. A number of milliseconds would drop the
microseconds of the default TIMESTAMP(6). Settings that would otherwise fail
silently are rejected at startup: a zero batch_size, which the client accepts
and which makes every poll come back empty, a negative starting offset, and only
one of the two SASL credentials.

Signed-off-by: seokjin0414 <sars21@hanmail.net>
Exercises the whole path rather than the connector in isolation: rows are
appended to a real Fluss log table, and the tests assert on what comes back out
of the Apache Iggy topic, including the offsets and bucket the connector
attached as metadata.

Two cases cover the ways a row could be skipped. One kills the Apache Iggy
server so a batch is rejected after the scanner has already handed its rows
over, then restarts it and expects every row, which only holds if the source
rewound the scanner. The other starts from `latest` over a table that already
holds rows, restarts the runtime before any new row arrives and expects the rows
written while it was down, which only holds if the tail resolved at startup
reached disk.

Another writes a row holding every supported column type, plus a row of nulls,
and checks the JSON each one becomes. Those rows are decoded from Arrow by a
real server, the path production reads, where the unit tests build rows in
memory.

The cluster runs as a single container. The image ships local-cluster.sh, which
starts the same embedded ZooKeeper, coordinator server and tablet server, but it
rewrites the tablet server's bind port to 0 so it cannot collide with the
coordinator, and a random port inside the container cannot be published to the
host. Giving the tablet server an explicit second port instead keeps both
reachable, and lets both servers advertise localhost, which then resolves the
same way inside the container and from the test process.

The table is created during fixture setup because the connectors runtime starts
before the test body and the connector resolves the table schema while opening.
Writes retry because bucket leadership is assigned shortly after the tablet
server registers, so the first attempts can still be rejected.

Signed-off-by: seokjin0414 <sars21@hanmail.net>
Registers the connector as a regular workspace member so it inherits the shared
dependency versions and stays inside cargo sort, the version bump script and the
DAG-based test scoping, the way every other connector does.

fluss-rs 1.0 ships its generated protobuf code, so building the workspace needs
no protoc. It does pull in arrow 59 while the workspace stays on 58 for iceberg
and deltalake, so the lockfile carries a second arrow line until those catch up.

Signed-off-by: seokjin0414 <sars21@hanmail.net>
A row the JSON mapping cannot hold used to fail its whole batch. The
failure repeats on every read and the SDK only logs a poll error, so
the source stopped making progress while still reporting itself as
running. Such a row is now dropped with an error naming its bucket and
offset, the offsets move past it, and close() logs how many rows were
skipped, as the runtime and the HTTP sink do with messages they cannot
decode.

Bad settings now fail before the connection is made: empty connection
settings, an empty or repeated column list, and a poll interval and
poll timeout that are both zero. With include_metadata on, a column
under the _fluss_ prefix is rejected instead of being overwritten.
Saved offsets for buckets the table no longer has are ignored, and the
projection is resolved once so the row decoder and the scanner agree on
column positions.

Signed-off-by: seokjin0414 <sars21@hanmail.net>
@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 new Apache Fluss source connector keeps the at-least-once invariant: offsets move only on ACK, a rejected batch rewinds to the committed offsets, and the resolved start offsets reach disk on the first acknowledged poll. The findings are documentation and robustness fixes: two README claims do not match what the runtime does with origin_timestamp and with the plugin config, the config struct stays serializable through an exposing secret helper, the skip branch classifies every mapping error as permanent, and the row mapper copies each column name; the reachability of a transient row-read failure cannot be settled from this checkout, so that item stays a warning rather than a blocking one.

Counts: critical 0, warning 5, nit 2, simplification 0


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/sources/fluss_source/src/lib.rs Outdated
Comment thread core/connectors/sources/fluss_source/src/lib.rs
Comment thread core/connectors/sources/fluss_source/src/mapping.rs Outdated
Comment thread core/connectors/sources/fluss_source/README.md Outdated
Comment thread core/connectors/sources/fluss_source/README.md Outdated
Comment thread core/connectors/sources/fluss_source/src/lib.rs Outdated
Comment thread core/integration/tests/connectors/fixtures/fluss/container.rs Outdated
@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
The config struct derived Serialize through the exposing secret helper,
so the annotation on sasl_password only looked like redaction. Nothing
serializes a plugin config, so the derive is gone and serializing the
password no longer compiles.

Each row was collected into a serde_json map first, which copied every
column name for every row. Rows are now written straight into the
payload from the schema and the Arrow batch. The JSON is unchanged,
checked byte for byte against the old mapping on random rows.

The skipped-row count was bumped while the batch was built, so a row
read again after a NACK was counted again. It now commits with the
offsets on ACK.

The README no longer says the runtime forwards origin_timestamp or that
the password is redacted everywhere. The test fixture keeps both port
listeners open until both ports are read, so the kernel cannot hand out
the same port twice.

Signed-off-by: seokjin0414 <sars21@hanmail.net>
@seokjin0414

Copy link
Copy Markdown
Contributor Author

Addressed the automated review in f8151c1. Six of the seven points are fixed. The transient column read cannot happen with fluss-rs 1.0.0, so I answered it in its thread and put the reason in the code.

/ready

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

hubcio commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

/skill team-review-slim

@github-actions

github-actions Bot commented Oct 2, 2026

Copy link
Copy Markdown

The review run ended with failure and produced no findings. Run log.

The test was written before stage_batch_state took a skipped count, and
the argument was later filled with 0, so it modelled three offsets
advancing with nothing produced and nothing skipped. Pass the three
skipped rows a real poll would report and check they reach the pending
state, since rows_skipped is not serialized.

Signed-off-by: seokjin0414 <sars21@hanmail.net>

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-review PR is waiting on a reviewer

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants