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]
