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]
