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]

Reply via email to