Thank you Anton, Lorenzo for the feedback.

The reviews have surfaced massive improvements opportunity that I missed
and eliminate unnecessary trade offs.

I’ll take this back to the drawing board and once ready follow up with a
response on this thread to all items brought up.

Best regards
Keith

On Tue, 23 Jun 2026 at 13:10, Anton Borisov <[email protected]> wrote:

> Thanks Keith,
>
> On zero-copy: agree the per-row re-encode is the wrong trade, and
> pushing materialisation to the client is a clean fix - it piggybacks
> the Arrow-to-row conversion clients already do, and since Flink, Spark
> and the lake reader all go through the Fluss Java client the stitch is
> largely write-once there, not per-connector. The one thing I'd keep in
> view is the splice I flagged first (keep base columns as the zero-copy
> file slices projection already serves, build only the enrichment
> column): its edge over client-side isn't CPU, it's that it changes no
> wire format and stays transparent to anything reading the fetch path,
> at the cost of a CRC pass server-side. Both are reasonable - worth
> picking on that wire-format/transparency question rather than on cost,
> since client-side's cost is smaller than it first looks.
>
> On the disk options, the thing that worries me is common to 2 and 3:
> the late enrichment, whenever it lands, still goes through the local
> appendColumns path, and that path can't write once the base for that
> range has left the local disk. Whether disk is freed by the timeout
> escape valve (2) or by tiering base ahead and truncating (3), the
> moment base is gone locally appendColumns reject offsets below
> localLogStart and enrichment for that range can never be written. The
> README's  retention note lands on exactly that: "deleting the local
> segment means enrichment data can never be written," with the
> fallbacks being hold-beyond-retention (disk pressure again) or accept
> enrichment data loss. So freeing the disk and still being able to
> enrich the freed range are in direct tension - the relief mechanisms
> work by deleting the base the enrichment needs.
>
> That's why I don't read option 2 as a disk fix, it relocates the
> problem - transient base-only or NULL columns in the lake, plus
> enrichment that may never land for timed-out ranges. Option 1 keeps
> completeness but takes the disk hit head-on.
> The way out, I think, is to stop gating disk on enrichment at all:
> tier base to the remote log at HW like a plain log table, so local
> disk is freed immediately. Because the base is then durably in remote
> storage, enrichment can be filled for that range at any time by
> reading the remote copy and writing a companion file - no local-log
> dependency, so no wedge, and it also removes the multi-job truncation
> worry, since local truncation then only depends on base tiering. The
> remote read/write plumbing already exists (RemoteLogDownloader on the
> read side, FsRemoteLogOutputStream and RemoteLogManifest on the write
> side), the genuinely new work is the control plane to write enrichment
> for an already-tiered range and track its completeness per-range in
> the manifest, rather than through the local-bound CEW.
>
> Given that, you can actually deliver the enriched lake the FIP is
> aiming at, without the disk trap. With base durable in remote, the
> lake materialises a range only once its enrichment has landed - the
> same backfill machinery, with the lake write deferred until the range
> is complete - so the lake ends up with every column, complete but
> lagging, never partial. That's the real target, and notably it does
> not need the Iceberg column-update, that primitive is only for making
> the enriched lake fresh rather than lagging, a later optimisation, not
> a gate on getting enriched data into the lake at all.
> For a table whose consumers are purely streaming and don't need
> enriched data at rest, base-only tiering plus live enriched reads from
> Fluss is the disk-safe minimum and a fine first increment - but it's
> the floor, not the destination, since on its own it doesn't give you
> the enriched lake. Option 2's timeout-fills-nulls is reasonable as an
> explicit opt-in for teams who'd rather have fresh-but-partial lake
> data, but I wouldn't make it the default or the disk answer.
>
> One framing note, +1 to Lorenzo: PK tables already do partial-column
> enrichment in place by targetColumns, with natural per-key parallelism
> and none of the CEW, merge-on-read, tiering or disk machinery. So a
> lot of what we're discussing is the cost of doing enrichment on
> append/log tables specifically. The FIP would be stronger scoping
> itself there explicitly - the cases where there's no key to merge on
> and offset-addressed replay matters - rather than as a general
> enrichment mechanism, since the keyed cases already have a path that
> avoids all of this.
>
> -- Anton
>
> вт, 23 июн. 2026 г. в 11:54, Lorenzo Affetti via dev <[email protected]
> >:
> >
> > 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>
>

Reply via email to