Thanks for the proposal.

Overall, I think this is a very strong FIP that further aligns our API
capabilities with KafkaConsumer, particularly regarding multi-topic
consumption.

However, I have several concerns regarding the current API design:

1. The Connection serves as the entry point for the Fluss API,
providing interfaces for Admin (DDL) and Table (per-table read/write).
However, LogScannerGroup does not appear to be an object at the same
hierarchical level as Admin or Table.

2. I have frequently encountered user requests in the community for a
single API to support multi-table writes in Fluss. I believe we should
consider the design of multi-table write APIs alongside multi-table
read APIs to avoid potential redesigns in the future.

3. The use cases for multi-table read/write operations differ
significantly from those for single-table operations. I suggest
isolating these APIs to allow them to evolve independently rather than
making them interdependent. For instance, query push-down is currently
less critical for multi-table scenarios, and having LogScannerGroup
depend on LogScanner feels awkward from a usability perspective.

4. We should minimize changes to foundational APIs such as LogRecord
and ScanRecord to avoid significant impacts on existing pipelines,
particularly concerning performance and compatibility.

Therefore, I propose introducing a MultiTable interface as the entry
point for multi-table read/write operations, corresponding to the
single-table Table interface. Additionally, we can provide a
MultiTableLogScanner to support core multi-table subscription
functionality. I have drafted a design outline [1] for your reference.

What do you think?

Best,
Jark

[1] https://gist.github.com/wuchong/d4bebe0f3fe610f1de730a05ab520c7d

On Fri, 8 May 2026 at 16:54, Hongshun Wang <[email protected]> wrote:
>
> Hi devs,
>
> Recently, I have been working on the Fluss YAML Source to support
> full-database synchronization (whole-database CDC). During the
> implementation, I ran into two fundamental limitations in the
> current LogScanner API:
>
> 1. LogScanner can only subscribe to a single table. When a business needs
> to consume change logs from multiple tables simultaneously, users must
> create a separate LogScanner for each table and manually manage
> multiplexing at the application layer. This leads to thread proliferation,
> resource waste (duplicate connections to the same TabletServer), and the
> inability to unify consumption into a single poll call.
>
> 2. LogScanner truncates logs by expected schema, losing schemaId and raw
> data. The current implementation projects records according to the schema
> provided at scanner creation time, which strips the schemaId metadata and
> discards columns added by subsequent schema changes. For CDC scenarios, we
> need to faithfully forward the original record data along with its schema
> version to downstream systems.
>
> To address these issues, I'd like to propose FIP-43: Multi-Table
> Subscription API[1] to solve both problems.
>
>  Looking forward to your feedback and suggestions!
>
>
> Best,
>
> Hongshun
>
> [1]
> https://cwiki.apache.org/confluence/display/FLUSS/FIP-43++Fluss+Multi-Table+Subscription+API

Reply via email to