nleigh opened a new pull request, #10341:
URL: https://github.com/apache/paimon/pull/10341

   ### Purpose
   
   Supersedes #10340 following a source-branch rename to 
`u/nleigh/changelog-as-append`; the implementation commit is unchanged.
   
   Some Flink SQL pipelines consume emitted changelog payloads as independent 
events. An aggregate update emits an old and a new payload, and downstream 
consumers retain and count both positively. For example, `+I(value=0), 
-U(value=0), +U(value=1000)` must become three stored append events, including 
two copies of value 0.
   
   Paimon's ordinary append writer rejects retract records, filtering retracts 
loses required payloads, and a keyed table maintains relational state. 
Migrating a message-based pipeline with this event-retention requirement 
currently needs a custom conversion operator.
   
   This draft proposes `sink.changelog-as-append`, a default-off Flink SQL sink 
mode that materializes each received INSERT, UPDATE_BEFORE, UPDATE_AFTER and 
DELETE payload as an independent INSERT event.
   
   **This is a draft for design discussion.** The interface and implementation 
location are proposed; upstream agreement and issue assignment are pending. A 
feature request is drafted but has not been filed. Feedback is welcome on the 
SQL sink approach, physical fields versus writable metadata, and whether an 
existing audit-log/changelog configuration can provide equivalent retention and 
conversion-time semantics.
   
   ### Changes
   
   - Request the full changelog, including before-images and full deletes where 
available.
   - Deep-copy each payload before the existing append writer, preserving 
duplicates and protecting against object reuse. Populate configured nullable 
physical STRING/BIGINT fields with the original row kind and conversion-time 
epoch milliseconds.
   - Validate the mode for non-keyed file-store tables. Reject format tables, 
overwrite, `ignore-delete`, `rowkind.field`, key-only delete configuration and 
generated fields used as partition keys.
   - Document the SQL example, restrictions and timestamp/replay contract in 
the append streaming guide and generated connector option reference.
   
   ### Scope and tradeoffs
   
   The mode covers the Flink SQL table sink. DELETE and UPDATE_BEFORE become 
recorded data and do not remove earlier events; readers see INSERT records and 
can inspect the original kind field. The option is disabled by default, and the 
existing core append-writer guards remain in place.
   
   Conversion time is neither source event time nor checkpoint commit time. 
Replay can assign new timestamps and move records to different downstream 
windows; timestamps are not globally ordered across subtasks. The sink 
preserves records delivered to it and cannot reconstruct history omitted 
upstream. Retaining all changes increases stored record volume compared with 
keyed state.
   
   ### Tests
   
   Local verification at `514da05`, based on upstream `master` at `df43423`, 
using Java 11:
   
   | Flink version / profile | Result |
   | --- | --- |
   | 1.20.1 / `flink1` | 21 passed; zero failures, errors or skips, after clean 
recompilation |
   | 2.2.0 / `flink2` | 21 passed; zero failures, errors or skips |
   
   Each run covered `ChangelogAsAppendTest` (5), `ChangelogAsAppendSinkTest` 
(6), `ChangelogAsAppendITCase` (2), `ChangelogModeTest` (6), and the existing 
append-writer DELETE/UPDATE_BEFORE rejection cases (2). There were 42 test 
executions total.
   
   Coverage includes:
   
   - Converter: all four kinds, equal duplicates, null payloads, nested object 
reuse, generated-field validation, timestamp bounds and input immutability.
   - Sink: full changelog negotiation, unchanged defaults, incompatible 
tables/options, generated partition fields and overwrite validation after 
abilities are applied.
   - SQL integration: explicit four-kind input and a real aggregate generating 
before-images, after-images and a final delete. Tests read committed storage 
and assert exact payload counts, duplicate retention, original kinds, INSERT 
read semantics and conversion-time bounds.
   - Regressions: default changelog negotiation and append-writer retract 
rejection.
   
   The connector/documentation module dependency closure built successfully 
with `-Pflink2,spark3 -pl paimon-flink/paimon-flink-common,paimon-docs -am 
install -DskipTests`. Checkstyle, Spotless, Enforcer and license checks were 
enabled in local Maven build/test invocations. Generated connector 
documentation succeeded with `-Pflink2,spark3,generate-docs -pl paimon-docs 
package -DskipTests`; `git diff --check` passed.
   
   [Upstream documentation CI 
passed](https://github.com/apache/paimon/actions/runs/37001920496/job/110821387160).
 The local site dependency install failed DNS resolution for 
`registry.yarnpkg.com`. The linked documentation result is from #10340 at the 
same implementation commit; CI for this replacement PR is pending.
   
   The PR remains in draft while design discussion and the following validation 
continue:
   
   - Full supported Flink/JDK matrix, including Java 8, and remaining upstream 
CI.
   - Streaming reads, compaction retention and checkpoint/recovery/replay 
coverage for the agreed contract.
   - Application-specific end-to-end parity and adoption checks.
   


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