Thanks Hongshun for the updates. The current design looks good to me, and I agree with Anton's points. Once his comments are addressed, I think we can start the vote.
Best, Jark On Thu, 21 May 2026 at 19:13, Anton Borisov <[email protected]> wrote: > > Hi Hongshun, > > Thanks for iterating on this. I like the current shape(+1) and think > it's a very much needed addition. > I studied the FIP and have some thoughts, hope you find it useful. > > 1. Historical schemaId semantics on writes > The FIP specifies fail-fast behavior when a record's schemaId is newer > than the server's current schema, but does not specify what happens > when the record's schemaId is older. > This matters for the CDC use case in section 2. For example, a > pipeline may read a record at schemaId = 3 and write it back after the > target table has moved to schemaId = 5. > > If this is out of scope for phase 1, for example because it requires > partial-update / target-columns semantics, please state that > explicitly. Otherwise, please specify the expected writer behavior for > historical schemaIds. > > 2. Phase-1 scope > From the PoC, phase 1 appears to exclude: > - partial-update / target-columns on writes, implying MergeMode.DEFAULT only; > - projection / filter on multi-table reads; > - bounded multi-table reads, since MultiTableBatchScanner is currently > only a placeholder. > > The FIP already refers to "Phase 1" in the WriteResult code block. > Please add a "Scope (Phase 1)" to make the tradeoffs and scope > explicit. > > 3. Cross-table fetch fairness > The FIP does not specify what fairness guarantees MultiTableLogScanner > provides when subscribed buckets span multiple tables on the same > leader. > Please clarify whether the intended contract is per-bucket fairness > only, per-table fairness, or something else, and whether there is any > related configuration. > > One follow-up item: > the FIP does not specify what happens if a table is dropped and > recreated under the same TablePath while a scanner is subscribed. > Since Fluss identifies tables by tableId rather than path, recreate > produces a new tableId under a stable path. It would be useful to > document the recovery behavior or list this as a known limitation. > > What do you think? > > -- Anton > > чт, 21 мая 2026 г. в 10:20, Hongshun Wang <[email protected]>: > > > > Hi devs, > > > > I have modified this FLIP[1] again: > > 1. Rename MultiTableScannerRecord -> MultiTableScanRecord. > > 2. remove the optional schema ID behavior and require users to always > > specify a schema ID when writing. > > > > Look forward to hearing from you. And I will start a vote soon if there is > > no problem any more. > > > > > > Best, > > Hongshun > > > > On Wed, May 20, 2026 at 11:17 AM Hongshun Wang <[email protected]> > > wrote: > > > > > 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 > > >> > >> > > >> > > > > >> > > >
