Hi Keith, I went through FIP-45 and I think it is an amazing piece of work — the CEW/watermark design in particular is elegant and well thought out. I have a few questions and suggestions I'd love to discuss.
First, a minor note on the motivation (§1): the problem statement focuses on storage cost, but I think the stronger argument is pipeline coupling — today, a slow enrichment job stalls all downstream consumers. Surfacing this angle would make the design choices, especially CEW gating, even more compelling. On the API side (§3), I wonder whether writing enrichment through a separate proxy sink table is the right ergonomics. The FIP promotes the idea of multiple independent jobs enriching different column groups of the same base table, but the DDL requires each job to declare its own virtual table with enrichment.target pointing back. A direct write to the base table with a job-level 'enrichment.group' option would feel more natural and would also let the connector handle _bucket/_offset resolution internally, avoiding leaking those internals into user SQL. Happy to be wrong here if there is a reason I'm missing. A concern I have on the write path (§2, §4): unlike standard Fluss log writes — which are naturally multi-writer thanks to per-writer writerId and batchSequence tracking — enrichment writes share a single global cursor per (bucket, group). It would be good to clarify how the Flink connector is expected to enforce single-writer-per-bucket routing, and who owns that responsibility. One broader question I think would strengthen the FIP: it would be valuable to include a short comparison with PK tables. For PK tables, the enrichment problem is already handled natively — partial column updates via targetColumns let a writer fill specific columns for an existing key, with the KV store merging them in-place, and parallelism is natural since writers operate on independent keys. The FIP could explicitly call out why this mechanism does not apply to log tables (no primary key, no update semantics, offset-addressed records) and how the proposed CEW/segment approach is the log-table equivalent. This would make the design feel more complete and help readers understand the scope of the proposal. Best, Lorenzo On Tue, Jun 23, 2026 at 12:38 AM Keith Lee <[email protected]> wrote: > Hello, > > Thank you Giannis and Anton for your feedback. The pushback on zero-copy is > valuable, the merge-on-read approach on the branch was built off the > motivation of making current APIs unchanged for polling logs, I am > convinced now that it was the wrong trade off to make and retaining > zero-copy can be done by pushing materialisation to client side (which > would be virtually free anyway as it is a conversion to columnar row that > clients already do.) > > > LogTieringTask is untouched, there's no escape valve or backfill path > > This is correct. I considered making an escape valve e.g. configurable > timeouts where timeouts result in null (we will have to force column type > to be nullable) or sentinel values. Both of these add complexities to the > design and neither null or sentinel values are easily distinguishable from > true nulls or sentinel equivalent values downstream. > > Your input would be very useful here as this aspect of the design is about > which trade off Fluss should make: > > 1. Preserve no partial row, simplicity in design and interpreting nulls: > this is the trade off that the design current takes. You have a point that > this story can be fleshed out further. I believe we can let users decide > which columns are necessary to be tiered, (as an example, they can decide > to just tier the base columns). If tiering is blocked by enrichment, the > failure mode is in a way similar to failing tiering job. > > 2. Timeout and fill with nulls/sentinel value: this is an alternative so > that tiering works and Fluss cluster storage do not get filled up if > enrichment jobs are slow or failing. This works but downstream consumers > will need to decide how to detect and handle nulls or sentinels from > timeouts and the backfill workflow will have to be orchestrated by user > outside of Fluss. > > 3. Same as 1. but allow for breaking up of tiering jobs and utilise tiered > storage feature to append enrichment columns. In this model, we would build > on 1. and also allow for multiple tiering jobs per log table for base and > enriched parts. This model would need to use currently unavailable feature > such as efficient column update on Iceberg [1]. There is added complexity > in this model in deciding when can log table be truncated. As tiering is > split up to multiple jobs, no single job can truncate rows after tiering as > other jobs might not have tiered up to that offset yet. > > > > > what does a lake reader see for those columns until backfill lands, given > the "no partial rows" contract? > > As is, lake reader will not see those columns or rows as tiering is blocked > by enrichment, fulfilling the no partial rows contract. However, if we make > trade off #2 mentioned earlier, they would either see nulls or sentinel > values. > > I’d appreciate your thoughts on the above before I address the other > comments or updating the design and draft implementation. > > Best regards > Keith > > > [1] > > https://docs.google.com/document/d/1Bd7JVzgajA8-DozzeEE24mID_GLuz6iwj0g4TlcVJcs/mobilebasic > > On Mon, 22 Jun 2026 at 04:22, Anton Borisov <[email protected]> wrote: > > > Hi Keith, > > > > Thanks for the writeup, I went through the FIP and the branch. > > The CEW completeness contract is the right core idea. > > > > A few concerns, aside from what Giannis already called out about > zero-copy. > > Fully zero-copy isn't achievable here, but the merger re-encodes every > > column row by row when it could keep the base columns as the zero-copy > > file slices projection already serves and build only the enrichment > > column - so the per-row rebuild is the costly extreme, not the floor. > > > > Okay, back to the points: > > 1. Leader failover. The FIP says reads gate at CEW so a survivor > > "never reveals enrichment that wasn't durably replicated". But > > onBecomeNewLeader seeds CEW > > from the promoted replica's own local EWM, and enrichment lag doesn't > > gate ISR membership, so a replica that's in-sync on the base log but > > behind on enrichment can > > be cleanly promoted with CEW set above what the ISR actually holds. It > > can also go backwards: promote a leader at EWM 10 and CEW jumps to 10, > > it dies, an in-sync > > follower still at EWM 4 is promoted and CEW drops to 4 and rows > > clients already read as enriched are hidden again. > > testNewLeaderSeedsCewFromLocalEwm encodes that over-claim as the > > expected result. Could CEW be checkpointed the way the high watermark > > already > > is (checkpointHighWatermarks runs periodically, off the hot path) and > > seeded from that on promotion? A stale checkpoint only under-claims, > > which is safe. > > 2. On tiering and disk: I first thought a slow enricher could wedge a > > bucket permanently, but the companion-file backfill design handles > > that, so - fine. > > What's left is that none of it is in the branch yet - LogTieringTask > > is untouched, there's no escape valve or backfill path, and > > computeTierSafeEndOffset only caps the > > lake split. So the disk-overflow story is unbuilt and untested in the > > branch, and it's the heaviest part of the proposal (companion remote > > segments, manifest rewrites, a second merge on the lake side). > > One question while it's still on design paper: during the escape-valve > > window, segments tier before enrichment is done - what does a lake > > reader see for those columns until backfill lands, given the "no > > partial rows" contract? > > > > A few smaller things: > > - commitTimestamp looks like it gets dropped in the merge, the Arrow > > builder takes no timestamp and setCommitTimestamp only runs on the > > leader's append path, so enriched reads > > and the lake __timestamp carry the builder default, not the > > original. baseLogOffset does survive. Worth confirming. > > - listOffsets(LATEST) returns CEW, not HW, for every caller on a > > column-group table - the tiering split's CEW bound > > (computeTierSafeEndOffset) leaks into the client's latest-offset > > answer, so a base-only bounded scan stops short of HW. > > Should listOffsets just return HW and keep the CEW bound to the > > tiering generator, since enrichment reads are already gated at fetch > > time? > > - Giannis, on offset 0 never being enrichable - I don't think it's > > real. The write check is source_offset == EWM and EWM starts at 0, so > > offset 0 goes through; it's > > the javadoc ("starts at -1") and the error string that misleads. > > The convention already matches the high watermark. > > > > Let me know what you think. > > > > -- Anton > > > > пт, 19 июн. 2026 г. в 06:23, Giannis Polyzos <[email protected]>: > > > > > > Thanks Keith, solid proposal and the compat/CEW design is great. > However > > I > > > have some concerns/blockers. > > > > > > Two blockers: 1. merge-on-read drops the zero-copy fetch path that > > > projection keeps today via FileLogProjection; the per-row re-encode > hits > > > every enrichment projection and all SELECT *. > > > 2. the re-encode threads no batch metadata (commitTimestamp, > > baseLogOffset, > > > lastOffsetDelta), corrupting offset/time-travel reads. > > > > > > Two things that could be addressed: The OffsetIndex is sparse and > > > slot-indexed, but EnrichmentSegment appends per-row and treats the > offset > > > as the ordinal, overflowing past 2^31 with an O(batch²) walk. EWM is > > > inconsistent (starts at 0, gate min(HWM,EWM), write at EWM+1), so > offset > > 0 > > > is never enrichable; Can we use the exclusive HWM convention. > > > > > > Also, tier-eligibility needs all groups caught up, so a stalled job > pins > > > segments and fills disk. > > > > > > Let me know your thoughts. > > > > > > Best, > > > Giannis > > > > > > > > > On Tue, 16 Jun 2026 at 12:19 AM, Keith Lee <[email protected]> wrote: > > > > > > > Hello dev, > > > > > > > > I have also updated the motivation section to include the convergence > > of > > > > lakehouse table formats that we are seeing in this area. > > > > > > > > Would appreciate further feedback from the community. > > > > > > > > Best regards > > > > Keith > > > > > > > > On Mon, Jun 15, 2026 at 10:17 PM Keith Lee <[email protected]> > wrote: > > > > > > > > > Hello Zhe, > > > > > > > > > > Thank you for your review, I have incorporated changes based on > your > > > > > feedback; I appreciate it as they make the proposal more robust. > > > > > > > > > > Best regards > > > > > Keith > > > > > > > > > > On Fri, Jun 12, 2026 at 4:29 AM Zhe Wang <[email protected]> > > wrote: > > > > > > > > > >> Hi Keith, > > > > >> > > > > >> Thanks for the detailed FIP. I read through the proposal and the > > overall > > > > >> direction makes sense to me. I have two questions/suggestions > about > > the > > > > >> Flink SQL user-facing part and the column-group metadata. > > > > >> > > > > >> 1. Literal-only enrichment SELECT > > > > >> > > > > >> The caveat around literal-only enrichment expressions looks a bit > > > > >> surprising to users: > > > > >> > > > > >> INSERT INTO device_logs_geo_sink > > > > >> SELECT _partition, _bucket, _offset, 'London' > > > > >> FROM device_logs; > > > > >> > > > > >> According to the FIP, this can silently emit zero rows because no > > base > > > > >> column is referenced, so the connector may fall back to a > projection > > > > that > > > > >> is clamped by the enrichment gate. > > > > >> > > > > >> Would it be possible to make this fail at planning time with a > clear > > > > >> validation error, or have the source side automatically include a > > > > minimal > > > > >> base-column projection when metadata columns are requested for > > > > enrichment? > > > > >> Silent zero-output behavior feels hard to diagnose in production, > > > > >> especially because the SQL statement looks valid and deterministic > > from > > > > >> the > > > > >> user's point of view. > > > > >> > > > > >> 2. Column group name validation > > > > >> > > > > >> Column group names appear in multiple places: > > > > >> > > > > >> - DDL option keys, e.g. `column-groups.enriched_geo` > > > > >> - Java API values, e.g. `columnGroup("enriched_geo", ...)` > > > > >> - RPC fields > > > > >> - physical segment files, e.g. `.col.<group>.log` > > > > >> > > > > >> Should the FIP explicitly define the allowed format for column > group > > > > >> names? > > > > >> For example, whether names may contain `.`, `/`, spaces, uppercase > > > > >> characters, or characters that need escaping. > > > > >> > > > > >> A small validation rule such as `[A-Za-z_][A-Za-z0-9_]*` or > another > > > > >> clearly > > > > >> documented convention would make DDL round-tripping, file naming, > > and > > > > >> error > > > > >> messages easier to reason about. > > > > >> > > > > >> Best regards, > > > > >> Zhe Wang > > > > >> > > > > >> > > > > >> Keith Lee <[email protected]> 于2026年5月25日周一 22:49写道: > > > > >> > > > > >> > Hello devs, > > > > >> > > > > > >> > I have raised a proposal on Log Enrichment via Append Columns > [1]. > > > > >> > > > > > >> > A common Fluss enrichment pipeline today looks like the > following: > > > > >> > > > > > >> > edge device --> log table A --> Flink job (enrichment) --> log > > table B > > > > >> --> > > > > >> > consumer > > > > >> > > > > > >> > If A has 50 columns and the enrichment adds 3, the 50 base > > columns are > > > > >> > stored twice - once in A, once in B. At scale this is real > > operational > > > > >> and > > > > >> > infrastructure cost - 2x storage, 2x write I/O, 2x fetch I/O for > > any > > > > >> > downstream that consumes B plus the additional architectural and > > > > >> > operational complexity of running and maintaining a two > pipelines. > > > > >> > > > > > >> > We propose a first-class primitive that lets the same row gain > > columns > > > > >> over > > > > >> > time, with the system enforcing "completeness" at the read > > boundary. > > > > >> With > > > > >> > this proposal, enrichment job(s) writes enrichment columns back > > to the > > > > >> same > > > > >> > log table A via appendColumns; there is no log table B. > Different > > > > >> > enrichment tasks can also happen at different cadence e.g. > > streaming > > > > for > > > > >> > GeoIP lookup, batch for vector embedding. > > > > >> > > > > > >> > Base columns are stored once and enrichment columns are stored > > once > > > > (per > > > > >> > column group). Consumers project whatever they need against the > > single > > > > >> > table: projections that touch only base columns advance up to > high > > > > >> > watermark (HWM) as today; projections that touch enrichment > column > > > > >> group g > > > > >> > are clamped by enrichment water mark (EWM) at min(HWM, EWM_g), > so > > > > >> partial > > > > >> > rows are never observed and the system not the consumer owns the > > "is > > > > >> > enrichment caught up?" decision. > > > > >> > > > > > >> > There are more details on the FIP and POC implementation [2]. > > > > >> > Looking forward to your comments > > > > >> > > > > > >> > Best > > > > >> > Keith > > > > >> > > > > > >> > References > > > > >> > [1] > > > > >> > > > > > >> > > > > > >> > > > > > > > https://cwiki.apache.org/confluence/display/FLUSS/FIP-45%3A+Log+Enrichment+via+Append+Columns > > > > >> > [2] > > > > >> > > > > > >> > > > > > > > https://github.com/leekeiabstraction/fluss/tree/option02-lateMaterialized > > > > >> > > > > > >> > > > > > > > > > > > > -- Lorenzo Affetti Senior Software Engineer @ Flink Team Ververica <http://www.ververica.com>
