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 > > >
