nleigh opened a new issue, #10343: URL: https://github.com/apache/paimon/issues/10343
### Search before asking - [x] I searched in the [issues](https://github.com/apache/paimon/issues) and found nothing similar. ### Motivation Some Flink SQL pipelines consume emitted changelog payloads as independent events. Downstream consumers retain and count each payload positively, including before-images and deletes. Migrating such a pipeline to Paimon requires retaining these payloads and exposing when each reaches the conversion operator so downstream consumers can assign it to a time window. For example, an aggregate can emit: <br class="Apple-interchange-newline"> Received change | Required stored event -- | -- INSERT: value = 0 | INSERT: value = 0, original kind = INSERT UPDATE_BEFORE: value = 0 | INSERT: value = 0, original kind = UPDATE_BEFORE UPDATE_AFTER: value = 1000 | INSERT: value = 1000, original kind = UPDATE_AFTER The result must contain three stored records, including two copies of value 0. A DELETE must similarly become a stored event containing its received payload. Ordinary append writing rejects retract records, filtering retracts loses required events, and keyed-table state does not itself provide this event-retention contract. We currently use a custom conversion operator and would like a reusable capability in the Paimon Flink SQL sink. ### Solution Add an explicit, default-off Flink SQL sink option, provisionally named sink.changelog-as-append. When enabled, the sink requests before-images and full delete payloads where available, deep-copies each received row, fills configured fields with the original row kind and conversion-time epoch milliseconds, and sends the row to the existing append writer as INSERT. Duplicate payloads are retained. The proposed interface uses two nullable physical columns, STRING and BIGINT. The INSERT supplies typed NULLs that conversion replaces: ``` CREATE TABLE change_events ( entity_id STRING, total BIGINT, original_kind STRING, emitted_ms BIGINT ) WITH ( 'bucket' = '-1', 'sink.changelog-as-append' = 'true', 'sink.changelog-as-append.kind-field' = 'original_kind', 'sink.changelog-as-append.time-field' = 'emitted_ms' ); INSERT INTO change_events SELECT entity_id, SUM(amount), CAST(NULL AS STRING), CAST(NULL AS BIGINT) FROM source_updates GROUP BY entity_id; ``` Initial scope is the Flink SQL sink for non-keyed file-store tables. The draft rejects format tables, overwrite, ignore-delete, rowkind.field, key-only delete configuration, and generated fields used as partition keys. Format tables are rejected when their SQL sink is created, including for batch inserts. The conversion operator preserves upstream parallelism and its configured flag, avoiding an extra rebalance or additional writers after a global aggregate. The option is disabled by default; normal sink behavior and core append-writer guards are preserved. Each record delivered to the sink becomes an event; upstream history that was omitted or combined cannot be reconstructed. UPDATE_BEFORE and DELETE are stored data and do not retract earlier events. Conversion time is neither source event time nor checkpoint commit time: replay can assign a new timestamp and change downstream window assignment. Timestamps are not globally ordered across workers. Retaining all changes increases storage volume. ### Anything else? A draft implementation is available in [PR #10341](https://github.com/apache/paimon/pull/10341). The interface and implementation location remain proposed; upstream design agreement and issue assignment are pending. Feedback requested: 1. Does this belong in the SQL sink or a dedicated event-oriented API/operator? 2. Should generated kind/time values use physical columns or writable metadata, and should both be mandatory? 3. Can an existing audit_log/changelog configuration provide equivalent durable event retention and conversion-time semantics? 4. What restrictions are needed for option changes, schema evolution and savepoint restores? I searched existing issues and found no equivalent request. The closest was this [#2083](https://github.com/apache/paimon/issues/2083) — Retain CDC deletes and operation metadata, which has some overlap. I am willing to contribute the implementation. ### Are you willing to submit a PR? - [x] I'm willing to submit a PR! -- 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]
