201811510411lw commented on issue #12243: URL: https://github.com/apache/seatunnel/issues/12243#issuecomment-5628611675
@DanielLeens Thanks for requesting upstream validation. I have now reproduced this on upstream dev commit `a0561d4045281c8e31ca96c6812b5024534db481`, and submitted the regression as draft PR https://github.com/apache/seatunnel/pull/12264 for independent reproduction. Only the test and test-scoped dependencies were added; no production source code was changed, and no downstream routing code or custom sink subclass was used. **Environment and effective configuration** - Java 8; Flink 1.18.1 with the SeaTunnel Flink 1.15 adapter/starter; Paimon 1.1.1. - Temporary local filesystem warehouse; `bucket=1`, `write-only=true`. - Primary key `(id, part)`; no partition keys. - Two upstream tasks; sink parallelism 2, with parallelism 1 as the control. - 100 INSERTs, DELETE of `(99, part-1)`, and four UPDATE_AFTER events for `(98, part-0)` ending in `after-3`. **Results through the actual upstream Flink sink integration** - Two writers: all three repetitions failed. A fresh Paimon read returned 100 rows instead of 99; the deleted key remained and the updated key retained its old value. - Two writers, emitting changes only after the INSERT checkpoint completed: all three repetitions failed identically. - Single-writer controls: both scenarios passed, including comparison of all complete primary keys and values. The test PR contains the runnable reproducer. From the PR branch, with Java 8 selected: ```bash mvn -B -pl seatunnel-connectors-v2/connector-paimon -am \ -Dskip.spotless=true \ -Dtest=PaimonUpstreamDeleteTest \ -Dsurefire.failIfNoSpecifiedTests=false verify ``` Expected result on this revision: **8 tests, 6 assertion failures, 0 errors, and 2 passing single-writer controls**. The PR intentionally remains a reproduction-only draft; it does not contain a fix. Separate build verification with `-DskipTests verify` passed, and runtime class origins confirmed that the test loaded the isolated upstream build artifacts. Regarding `PaimonBucketAssigner`, I found a distinction relevant to this case: [the current writer enables it only for `HASH_DYNAMIC`](https://github.com/apache/seatunnel/blob/a0561d4045281c8e31ca96c6812b5024534db481/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/sink/PaimonSinkWriter.java#L174). This fixed-bucket case takes [the `tableWrite.write(rowData)` branch](https://github.com/apache/seatunnel/blob/a0561d4045281c8e31ca96c6812b5024534db481/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/sink/PaimonSinkWriter.java#L230), with an independent table writer per sink subtask. [The Flink starter path](https://github.com/apache/seatunnel/blob/a0561d4045281c8e31ca96c6812b5024534db481/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/SinkExecuteProcessor.java#L51) does not establish a single writer owner for that fixed bucket. Could you please review the reproducer in #12264 and confirm whether this is the expected upstream regression coverage before we discuss a fix? -- 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]
