goutamadwant opened a new issue, #12082:
URL: https://github.com/apache/seatunnel/issues/12082

   ## Search before asking
   
   - [x] I searched the open and closed issues and pull requests using 
`AmazonSqs`, `canal_json`, `debezium_json`, `DeserializationSchema`, and 
`Collector`. I found no issue or active PR covering this defect.
   
   ## What happened
   
   The Amazon SQS source documents `canal_json` and `debezium_json` as 
supported `format` values, but both formats fail before producing any rows.
   
   `AmazonSqsSourceFactory` correctly creates `CanalJsonDeserializationSchema` 
or `DebeziumJsonDeserializationSchema`. However, `AmazonSqsDeserializer` always 
invokes `DeserializationSchema.deserialize(byte[])`. Both CDC schemas reject 
that single-row overload and require `deserialize(byte[], 
Collector<SeaTunnelRow>)`, because one CDC message can produce zero, one, or 
multiple rows.
   
   For example, a Canal or Debezium update should emit `UPDATE_BEFORE` and 
`UPDATE_AFTER`. Instead, the first SQS message fails with:
   
   ```text
   java.lang.UnsupportedOperationException: Please invoke 
DeserializationSchema#deserialize(byte[], Collector<SeaTunnelRow>) instead.
   ```
   
   This makes two of the four documented Amazon SQS formats unusable.
   
   ### Reproduction
   
   1. Configure an Amazon SQS source with `format = canal_json` and enqueue a 
valid Canal update message:
   
   ```json
   
{"data":[{"value":"after"}],"old":[{"value":"before"}],"database":"inventory","table":"orders","type":"UPDATE"}
   ```
   
   2. Run the source. The same failure is reproducible with `format = 
debezium_json` and a valid Debezium update envelope.
   3. Observe that the poll throws the `UnsupportedOperationException` above 
before either CDC row is emitted.
   
   I also reproduced both paths with factory-created source readers in focused 
connector tests: Canal on Java 8 and Debezium on Java 11.
   
   ### Proposed contract
   
   - Use the collector-based deserialization overload for `canal_json` and 
`debezium_json`.
   - Preserve zero-row CDC events such as Canal query/DDL filtering instead of 
treating them as failures.
   - Preserve row order for multi-row events.
   - Buffer the rows from one SQS message before sending them to the real 
downstream collector. This prevents tolerant format parsing from swallowing 
downstream collector failures.
   - Delete an SQS message only after every row produced from that message has 
been collected successfully.
   - With `ignore_parse_errors = false`, a malformed CDC message must fail the 
poll and remain available for redelivery.
   - With `ignore_parse_errors = true`, a malformed CDC message must be skipped 
atomically. No rows from a partially parsed message should be emitted; deletion 
continues to follow `delete_message`.
   - Preserve the existing `json` and `text` behavior and public 
constructor/interface contracts.
   
   ## SeaTunnel Version
   
   Current `dev` at commit `98182ca59c93c10390717af4a49b8806f5c555f1`.
   
   ## SeaTunnel Config
   
   ```conf
   source {
     AmazonSqs {
       url = "http://sqs-host:4566/000000000000/orders";
       region = "us-east-1"
       access_key_id = "test"
       secret_access_key = "test"
       format = canal_json
       delete_message = true
       ignore_parse_errors = false
       schema = {
         fields {
           value = string
         }
       }
     }
   }
   ```
   
   Changing `format` to `debezium_json` reproduces the same 
unsupported-overload failure with a valid Debezium envelope.
   
   ## Running Command
   
   ```shell
   bin/seatunnel.sh --config config/amazonsqs-cdc.conf -e local
   ```
   
   The deterministic connector-level reproduction uses factory-created readers 
and does not require a live AWS account.
   
   ## Error Exception
   
   ```log
   java.lang.UnsupportedOperationException: Please invoke 
DeserializationSchema#deserialize(byte[], Collector<SeaTunnelRow>) instead.
   ```
   
   ## Zeta or Flink or Spark Version
   
   Connector-level defect; it occurs before engine-specific processing.
   
   ## Java or Scala Version
   
   Reproduced on Java 8 and Java 11.
   
   ## Are you willing to submit PR?
   
   - [x] Yes, I am willing to submit a PR.
   
   ## Code of Conduct
   
   - [x] I agree to follow this project's Code of 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