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]

Reply via email to