201811510411lw opened a new issue, #12243:
URL: https://github.com/apache/seatunnel/issues/12243

   ### Search before asking
   
   - [x] I had searched in the 
[issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22)
 and found no similar issues.
   
   
   ### What happened
   
   ### Problem
   
   When writing CDC data to a Paimon primary-key table through SeaTunnel on 
Flink with multiple sink writers, deleted records may remain in Paimon after 
the writes are committed.
   
   Expected: a committed DELETE removes the corresponding primary key.
   Actual: the deleted record is still returned by a fresh Paimon read.
   
   ### Reproduction evidence
   
   A downstream Flink MiniCluster regression uses two source subtasks, two sink 
writers, a fixed-bucket Paimon table (`bucket = 1`), and synthetic 
INSERT/UPDATE/DELETE events.
   
   After inserting 100 distinct records and deleting one, the legacy sink path 
without bucket-aware routing fails the read-back assertion:
   
       expected: <99> but was: <100>
   
   The test exercises the actual SeaTunnel Flink sink path and Paimon connector 
with a local filesystem warehouse.
   
   This was reproduced on a customized SeaTunnel 2.3.13-based branch, not a 
clean upstream dev checkout. The regression test is currently downstream-only.
   
   ### Suspected cause
   
   The legacy Flink sink path appears not to guarantee a single writer owner 
for each physical Paimon table/partition/bucket. An INSERT and a subsequent 
DELETE for the same key may therefore reach different writers targeting the 
same bucket.
   
   The corresponding downstream regression passes after adding bucket-aware 
routing before the sink. This supports the diagnosis, but is not yet an 
upstream-verified fix.
   
   ### Related issues
   
   I checked #7769, #8565, #9595, #10978 and #11077. They discuss single-writer 
configurations, dynamic buckets or sharing writers across logical tables; I did 
not find an exact duplicate of this fixed-bucket DELETE case.
   
   Should the Flink integration enforce table/partition/bucket-aware routing 
before the Paimon sink? If this is already tracked elsewhere, please point me 
to the existing issue.
   
   ### SeaTunnel Version
   
   SeaTunnel 2.3.13-based customized downstream build.
   
   The failing reproduction uses a legacy-path control in the downstream 
regression suite.
   
   Reproduction on a clean upstream release or the latest upstream dev branch 
has not yet been performed.
   
   ### SeaTunnel Config
   
   ```conf
   The current reproducible case is a programmatic Flink MiniCluster regression,
   rather than a standalone SeaTunnel HOCON job.
   
   Test configuration:
   
   - Source parallelism: 2
   - Sink parallelism: 2
   - Connector: Paimon
   - Warehouse: temporary local filesystem directory
   - Database: routing_test
   - Table: table_0
   - Primary key: (id, part)
   - Paimon table options:
       bucket = 1
       write-only = true
   
   The test sends synthetic INSERT/UPDATE/DELETE records through the actual
   SeaTunnel Flink sink path, commits the writes, and reads the Paimon table 
back.
   
   The legacy control omits the downstream bucket-aware routing before the sink.
   
   No production credentials, endpoints or business data are required by this 
test.
   ```
   
   ### Running Command
   
   ```shell
   # Downstream regression control; this test is not available in upstream 
SeaTunnel.
   # Requires the downstream source tree and its built module dependencies.
   # The legacy control is expected to fail the DELETE correctness assertion.
   
   mvn -f seatunnel-connectors-v2/connector-paimon/pom.xml \
     -Dskip.spotless=true \
     -Dissue113.test.legacy-routing=true \
     
-Dtest=PaimonFlinkBucketRoutingTest#shouldRouteSingleBucketThroughActualFlinkSink
 \
     test
   ```
   
   ### Error Exception
   
   ```log
   The reported symptom is incorrect persisted data, rather than a connector 
exception.
   
   The downstream legacy-path regression failed with:
   
   expected: <99> but was: <100>
   
   This is a test assertion failure from reading the Paimon table after the 
write/delete sequence,
   not a production exception stack trace.
   ```
   
   ### Zeta or Flink or Spark Version
   
   Flink 1.18.1
   
   SeaTunnel Flink 1.15 adapter / starter integration.
   Paimon 1.1.1.
   
   ### Java or Scala Version
   
   Java 8
   
   ### Screenshots
   
   _No response_
   
   ### Are you willing to submit PR?
   
   - [ ] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://www.apache.org/foundation/policies/conduct)
   


-- 
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