goutamadwant opened a new pull request, #12084:
URL: https://github.com/apache/seatunnel/pull/12084

   ### Purpose of this pull request
   
   Closes #12082.
   
   The Amazon SQS source advertises `canal_json` and `debezium_json`, but its 
deserialization adapter always called the single-row `deserialize(byte[])` 
overload. Both CDC schemas reject that overload because one CDC message can 
produce zero, one, or multiple rows.
   
   This patch:
   
   - uses collector-based deserialization for Canal and Debezium JSON;
   - buffers all rows from one SQS message before emitting them downstream;
   - preserves valid zero-row CDC events;
   - emits update-before/update-after rows in their original order;
   - deletes an SQS message only after all buffered rows are collected 
successfully;
   - makes `ignore_parse_errors = true` skip a partially invalid CDC message 
atomically instead of emitting a partial result;
   - preserves the existing single-row interface method and constructors for 
compatibility; and
   - keeps the existing JSON and text paths unchanged.
   
   The per-message buffer is limited by the SQS message size and is used only 
for the two multi-row CDC formats. The normal JSON and text paths retain their 
existing single-row behavior.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes.
   
   Before this change, selecting `canal_json` or `debezium_json` caused the 
first message to fail with `UnsupportedOperationException`, so two documented 
formats were unusable.
   
   After this change, valid CDC messages produce all expected rows. For 
example, an update emits `UPDATE_BEFORE` followed by `UPDATE_AFTER`.
   
   With `delete_message = true`, the SQS message is deleted only after every 
emitted row is collected successfully. Parse-error behavior remains controlled 
by `ignore_parse_errors`.
   
   The English and Chinese Amazon SQS documentation now describe multi-row CDC 
emission and deletion ordering.
   
   ### How was this patch tested?
   
   Added focused factory-to-reader regression coverage for:
   
   - Canal update events producing two ordered rows;
   - Debezium update envelopes using the documented default schema-envelope 
setting;
   - valid Canal events that intentionally produce zero rows;
   - a downstream failure on the second row retaining the SQS message;
   - partially invalid multi-record Canal messages being skipped atomically 
when parse errors are ignored;
   - partially invalid messages failing without partial emission or deletion by 
default; and
   - unrelated runtime failures continuing to propagate.
   
   Validation results:
   
   - Java 8: the complete `connector-amazonsqs` module passed, 17/17 tests.
   - Java 11: the complete `connector-amazonsqs` module passed, 17/17 tests.
   - Java 11 module verification with tests skipped passed, including 
compilation, packaging, and Spotless.
   - `git diff --check` passed.
   
   The tests use a deterministic in-memory SQS client boundary and make no 
network calls.
   
   ### Check list
   
   * [x] No new Jar binary package or dependency is added.
   * [x] Updated the English and Chinese connector documentation for the 
user-facing behavior.
   * [x] No incompatible configuration or checkpoint-state change is introduced.
   * [x] Existing connector registration, plugin mapping, distribution POM, CI 
labels, E2E wiring, and plugin configuration remain unchanged because this 
patch fixes an existing connector implementation.


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