feat(connectors): add Apache Fluss source connector - #3799
seokjin0414 wants to merge 6 commits into
Conversation
|
Thanks for the PR. It is labeled Slash commands (own line, regular comment) move it around the queue:
See CONTRIBUTING.md for details. |
Codecov Report❌ Patch coverage is 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
🚀 New features to boost your workflow:
|
e8d3b85 to
67c05d3
Compare
67c05d3 to
26bf766
Compare
|
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 Thank you for your contribution! |
26bf766 to
d575c8f
Compare
d575c8f to
944561e
Compare
|
Would be good to convert this to a draft until upstream changes are merged. This will get autoclosed otherwise |
944561e to
dd68fbe
Compare
|
Rebased onto latest master and adapted to the source batch acknowledgments (#3855): offsets are now staged in 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. |
dd68fbe to
db9cfca
Compare
|
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 @hubcio this is the one the upstream change was for, ready for review when you have time. |
fa9c181 to
8847622
Compare
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>
053e3f2 to
828f4d4
Compare
|
/skill team-review-slim |
There was a problem hiding this comment.
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.
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>
|
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 |
|
/skill team-review-slim |
|
The review run ended with |
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>
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 ratherthan borrowing the connection.
Offsets are tracked per bucket and follow the batch acknowledgment contract
(#3855):
pollstages the candidate offsets with the batch andon_batch_resultcommits them onAck. The scanner moves past records as soonas it returns them, so a
Nackrewinds it to the committed offsets and thoserows 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_resultstops the source.A row the JSON mapping cannot hold, such as a date beyond what
chronorepresents, 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 thefirst row would resolve
latestagain and skip whatever was written inbetween. On restart, buckets already present in the state keep their offset
and only new ones fall back to
starting_offset, which keeps a widened bucketcount 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
columnsis 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_metadataon, thebucket, offset and timestamp are added under the
_fluss_prefix, so a tablecolumn 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:
TIMESTAMPwithout a timezone,TIMESTAMP_LTZas RFC 3339 in UTC. A number of milliseconds would drop themicroseconds of the default
TIMESTAMP(6).Settings that would otherwise fail late or silently are rejected before
connecting:
bootstrap_servers,databaseortablecolumnslistbatch_size, which the client accepts and which makes every poll comeback empty
starting_offsetpoll_intervaltogether with a zeropoll_timeout, which would pollan 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 plainworkspace member with no system dependency and no CI change, and the earlier
ci:commit is gone.One dependency note:
fluss-rs1.0 uses arrow 59 while the workspace stays on58 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:
offsets and bucket.
handed its rows over. Once the server is back every row still arrives, which
only holds if the source rewound the scanner.
latestover a table that already holds rows, the runtime isrestarted before any new row arrives. The rows written while it was down are
delivered and the older ones are not.
the expected JSON. These rows are decoded from Arrow by the server, the path
production reads, while the unit tests build rows in memory.
chronocannot represent sits between two good rows.The good rows arrive and the bad one is skipped.
Each value arrives under its own name and the third column is absent.
_fluss_, sowith
include_metadataon the source fails to start. The projection testreads 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 runfor 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 allpass.
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 thecontainer 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:
fluss-rs1.0 can read their changelog, but mapping itsinsert, update and delete records onto messages deserves its own PR.
payload_format = "arrow_ipc". It needs the batch scanner, which has its ownoffset-tracking path, so it is left for a follow-up.
The README also lists behavior a user should know about: rows that cannot be
converted are skipped, a crash right after starting from
latestcan missrows, and an offset outside what a bucket still holds is not reset.
Relates to #3688