JingsongLi commented on PR #9822: URL: https://github.com/apache/paimon/pull/9822#issuecomment-5833425147
I took another look at head `b0fd1ea`, focusing on the public semantics and schema evolution. The Cassandra `WRITETIME` use case is valid, and preserving the physical before-image is the right direction. I need to revise my earlier no-blocker assessment: I found two correctness issues with the current metadata surface. **1. [P1] Event metadata breaks retract semantics when used by normal SQL operators.** `LookupChangelogMergeFunctionWrapper` writes the *new event's* metadata into the `-U` record while its physical columns correctly contain the old row. For example, an initial `+I(event_ts=50, metadata=50)` followed by an update to 100 produces `-U(event_ts=50, metadata=100)`. In Flink, `WHERE metadata < 75` forwards the insert but drops its retraction, leaving a stale row. Grouping or aggregating by the metadata column has the same problem. The tests cover filtering on the physical `event_ts`, but not on the exposed metadata column. This is intrinsic to treating an event attribute as an ordinary column in a retracting dynamic table; please define an event-oriented consumption surface (or another mechanism that prevents relational operators from interpreting it as before-image data) and add a regression test for this sequence. **2. [P1] Changing `changelog-producer.metadata-field-prefix` loses historical metadata.** The prefix is not marked immutable, so it can be changed after changelog files exist. The write path stored the old prefixed field name, but `ChangelogEventMetadata.extraValueFields` and the reader construct only the *current* name for every file schema. After changing `__internal__` to `__event__`, older changelog files have no `__event__event_ts` field, so historical event metadata reads as `NULL`; the data-file fallback is intentionally disabled for changelog files. Please keep the on-disk identity stable across option changes, or reject prefix changes once files exist. A write → ALTER prefix → read-old-changelog test would make the contract clear. There is also a source/sink mismatch in the documentation example: `METADATA FROM '__internal__event_ts'` needs to be declared on the **Paimon source**. An external sink's `METADATA FROM` key belongs to that sink and works only if the sink advertises that writable key. Please show the source alias and an explicit mapping to the sink's actual timestamp input. Finally, the option names obscure the contract. `expose-field-as-metadata` persists a *list* of post-merge event values; it does more than expose a read-time field and does not always preserve the raw incoming value. A name such as `changelog-producer.event-metadata-fields` would be clearer. The configurable prefix currently couples the on-disk field name, Flink metadata key, and Spark-visible column name. I suggest separating the stable storage identity from engine-specific presentation. Given that the concrete use case is a Flink-to-Cassandra flow, I would also consider landing core/Flink support first and adding the Spark schema/write-path changes with a separate demonstrated consumer. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
