Hi jark, The design intention is to simplify usage. When no schema ID is specified, the writer defaults to using the table schema that was cached during the first subscribe call. During write time, it validates whether the number of columns in the row matches the cached schema.
The current "no schema ID" approach is admittedly coarse-grained — it only checks that the column count matches, but does not verify that the columns themselves are semantically correct. For example, if we later support ALTER COLUMN TYPE, a row could pass the column-count check but actually have mismatched column semantics, leading to silent data corruption. Given this, I think it's also reasonable to remove the optional schema ID behavior and require users to always specify a schema ID when writing. This would make the contract stricter and safer, especially as schema evolution capabilities grow. WDYT? Best, Hongshun On Tue, May 19, 2026 at 12:13 PM Jark Wu <[email protected]> wrote: > Hi Hongshun, > > The new FIP looks good to me in general. I only have some minor comments. > > 1. Rename MultiTableScannerRecord -> MultiTableScannRecord, > MultiTableScannerRecords -> MultiTableScanRecords to align with > ScanRecord. > > 2. MultiTableWriter supports MultiTableWriteRecord without schema ID. > What's the behavior when there is no schema id on the write record? > What happens when the write record has different schema? > > Best, > Jark > > On Mon, 18 May 2026 at 20:39, Hongshun Wang <[email protected]> > wrote: > > > > Hi Devs, > > > > I have modified this FLIP[1]: tables are registered dynamically on the > > first MultiTableLogScanner::subscribe(...) call rather than > > pre-defined in MultiTableScan:: > > withTable before MultiTableLogScanner is created. And also , the > > multiple table writer is added. > > > > Poc code is here: https://github.com/apache/fluss/pull/3140 > > > > > > [1] > > > https://cwiki.apache.org/confluence/display/FLUSS/FIP-43++Fluss+Multi-Table+Subscription+API > > > > Best, > > Hongshun > > > > On Thu, May 14, 2026 at 2:14 PM Hongshun Wang <[email protected]> > > wrote: > > > > > Hi Jark, > > > I agree with you, the new interface is better. The only issue is that > > > `MultiTableScan withTable(TablePath tablePath)` and `MultiTableScan > > > withTable(TablePath tablePath, int[] projectedColumns)` is unnecessary, > > > because multi-table subscription often involves dynamically adding new > > > tables to LogFetcher during runtime, rather than declaring them in > advance. > > > > > > Best, > > > Hongshun > > > > > > On Tue, May 12, 2026 at 7:21 PM Jark Wu <[email protected]> wrote: > > > > > >> 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 > > >> > > > >
