nleigh opened a new pull request, #10340: URL: https://github.com/apache/paimon/pull/10340
### Purpose 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. These mechanisms do not directly provide an append event log of every received payload, so migrating an existing message-based pipeline currently requires a custom conversion operator. This draft proposes `sink.changelog-as-append`, a default-off Flink SQL sink mode that copies every received INSERT, UPDATE_BEFORE, UPDATE_AFTER and DELETE payload into an independent INSERT before the existing writer. It preserves duplicate records and fills configured nullable physical STRING/BIGINT fields with the original row kind and conversion-time epoch milliseconds. DELETE becomes recorded data rather than removing an earlier event. The sink requests before-images and full deletes where available. The mode rejects primary keys, format tables, overwrite, `ignore-delete`, `rowkind.field`, key-only delete configuration and generated fields used as partition keys. Existing default behavior and core append writer guards are preserved. The streaming documentation includes a generic SQL example and the generated connector option reference is updated. **This is a draft for design discussion.** The interface and implementation location are proposed; upstream agreement and issue assignment are still pending. Feedback is welcome on whether this belongs in the SQL sink, whether generated values should use physical fields or writable metadata, and whether an existing audit-log/changelog approach can provide the required event retention and conversion-time contract. Conversion time is not source event time or checkpoint commit time. Reprocessing can assign new timestamps and move records to different downstream windows. There is no global ordering across subtasks. The mode preserves records delivered to the sink, not every original source event: it cannot reconstruct omitted before-images or history. Retaining all changes also increases stored record volume compared with keyed state. ### Tests Included coverage: - 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. - Existing default changelog negotiation and append-writer retract rejection are selected as regressions. Validation of this implementation (`514da05`), based on upstream `master` at `df43423`, with Java 11: - The Flink connector/documentation module dependency closure built and installed successfully with `-Pflink2,spark3 -pl paimon-flink/paimon-flink-common,paimon-docs -am install -DskipTests`. - **21 selected tests passed under each of `flink2` and `flink1` on Java 11 (42 executions total), with zero failures, errors or skips:** `ChangelogAsAppendTest` (5), `ChangelogAsAppendSinkTest` (6), `ChangelogAsAppendITCase` (2), `ChangelogModeTest` (6), and the existing append-writer DELETE/UPDATE_BEFORE rejection cases (2). - The Flink 1 run used `clean test` to recompile the connector under that profile. Java 8 and the full supported connector matrix remain to be verified. - Normal formatting, checkstyle, enforcer and license checks were enabled in the Maven build/test invocations. - Generated connector documentation succeeded with `-Pflink2,spark3,generate-docs -pl paimon-docs package -DskipTests`. `git diff --check` passed. - Documentation site validation is currently blocked because `yarn install --frozen-lockfile` cannot resolve `registry.yarnpkg.com` in this environment. The public dependency install was retried with network access and the failure persisted. Remaining before readiness/adoption: - Supported Flink/JDK matrix and upstream CI. - Streaming reads, compaction retention and checkpoint/recovery/replay coverage for the agreed contract. - Application-specific end-to-end parity and adoption checks. The implementation introduces no release/version bump or deployment change. -- 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]
