201811510411lw opened a new pull request, #12264:
URL: https://github.com/apache/seatunnel/pull/12264
### Purpose of this pull request
Related to #12243. This draft adds a failing regression for fixed-bucket
Paimon DELETE/UPDATE residue through the actual SeaTunnel Flink sink
integration. With two writers, inserting 100 keys and deleting one leaves 100
keys in a fresh Paimon read; the updated key also retains its old value. Both
single-writer controls pass.
The runtime source is unchanged from upstream dev commit
`a0561d4045281c8e31ca96c6812b5024534db481`. This is a reproduction-only draft,
not a fix or a merge-ready change: the new multi-writer tests are expected to
fail until the underlying issue is resolved.
The test uses the original `SinkExecuteProcessor`, `FlinkSink`, and
`PaimonSinkWriter`, without a custom sink subclass or downstream routing code.
Test-scoped starter/adapter and Flink clients/streaming dependencies are needed
to run that integration in a local MiniCluster. No production dependency or
runtime implementation is changed.
### Does this PR introduce _any_ user-facing change?
No. Only one test class and four test-scoped dependencies are added.
### How was this patch tested?
Environment: Java 8, Flink 1.18.1 with the SeaTunnel Flink 1.15
adapter/starter, and Paimon 1.1.1. The warehouse is a temporary local
filesystem directory; no external services or credentials are needed.
Effective table configuration:
- Table: `routing_test.table_0`; fields: `id INT NOT NULL`, `part STRING NOT
NULL`, `value STRING`.
- Primary key: `(id, part)`; partition keys: none.
- Options: `bucket=1`, `write-only=true` (plus the generated temporary table
path).
- Job mode: `STREAMING`; two upstream tasks; sink parallelism 2, or 1 for
controls.
- Input: INSERT 100 distinct keys; DELETE `(99, part-1)`; four UPDATE_AFTER
events for `(98, part-0)`, ending in `after-3`.
The finite-input case explicitly sends INSERTs through upstream task 0 and
changes through task 1. The checkpoint case uses two real parallel source
subtasks and sends changes only after the INSERT checkpoint completes (200 ms
checkpoint interval). After the job finishes, native Paimon read-back compares
all complete primary keys and values, with separate count, DELETE, and UPDATE
assertions.
| Scenario | Sink writers | Repetitions | Result |
| --- | --- | --- | --- |
| Finite input, delayed changes | 2 | 3 | All fail: 100 rows; deleted key
and old update value remain |
| Changes after completed INSERT checkpoint | 2 | 3 | All fail identically |
| Finite-input control | 1 | 1 | Pass: 99 rows, deleted key absent, updated
value correct |
| Checkpoint control | 1 | 1 | Pass, including full key/value comparison |
Local result: **8 tests, 6 assertion failures, 0 errors, 0 skipped**, in
27.974 seconds. Runtime CodeSource output confirms the relevant classes came
from the isolated upstream build.
From this 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 exit code: **1**, from the six data assertions above. Use `verify`
rather than `test` for this reactor invocation so the upstream shaded Config
dependency is packaged before its dependent module is compiled.
Independent build verification passed for all 27 selected reactor modules:
```bash
mvn -B -pl seatunnel-connectors-v2/connector-paimon -am \
-Dskip.spotless=true -DskipTests verify
```
The affected module's `spotless:apply` and subsequent `spotless:check`, plus
`git diff --check`, passed.
### Relevant current routing path
In the pinned upstream revision, `PaimonSinkWriter` sets `dynamicBucket`
only for `HASH_DYNAMIC` and initializes/uses `PaimonBucketAssigner` only in
that branch. Fixed buckets take the `tableWrite.write(rowData)` branch, and
each writer creates its own `paimonTable.newWrite(commitUser)`. The Flink
starter adds the schema handler and then calls `sinkTo`, without establishing a
unique writer owner for the fixed bucket.
- [Fixed/dynamic bucket
selection](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)
- [Fixed-bucket write
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)
- [Flink sink
integration](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)
### Check list
- No binary files or production dependencies are added.
- User documentation, incompatible-change notes, and new-connector
registration are not applicable to this reproduction-only draft.
- A local Flink integration regression and single-writer controls are
included; no production source changes are proposed.
--
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]