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]

Reply via email to