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

   ## Search before asking
   
   - [x] I searched issues and pull requests in all states and found no report 
or change that covers this behavior.
   
   ## What happened
   
   When `delete_message = true`, an Amazon SQS message that cannot be 
deserialized can be deleted from the queue.
   
   `AmazonSqsDeserializer` currently catches `IOException` and returns `null`. 
It also forwards a `null` result from the configured `DeserializationSchema`. 
`AmazonSqsSourceReader` passes that result to the collector and then deletes 
the SQS message after `collect` returns. On the Zeta source-collector path, 
`null` can be wrapped in a record and forwarded, so the SQS deletion can 
complete before a later failure exposes the invalid row.
   
   This bypasses the queue's visibility-timeout retry and dead-letter-queue 
redrive behavior. It also conflicts with the connector documentation, which 
says `delete_message` removes a message only after it is deserialized 
successfully.
   
   A focused regression against the production reader reproduced both failure 
forms:
   
   1. The configured deserializer throws `IOException`.
   2. The configured deserializer returns `null`.
   
   Before the fix, the deserialization failure was not propagated, the 
collector received `null`, and the invalid message's receipt handle reached 
`deleteMessage`.
   
   Expected behavior:
   
   - An `IOException` or `null` deserialization result must fail the current 
poll.
   - The failed message must not be collected or deleted.
   - Later messages in the same received batch must not be processed or deleted.
   - Messages successfully collected before the failed message may already have 
been deleted.
   - A collector failure must continue to prevent deletion.
   - Valid-message behavior and the default `delete_message = false` behavior 
must remain unchanged.
   
   ## Why the proposed fix helps
   
   The proposed fix converts both an `IOException` and a `null` result into a 
project-native unchecked connector exception before the source reader reaches 
`collect` or `deleteMessage`. This leaves the failed SQS message available 
after its visibility timeout and allows the queue's retry or dead-letter-queue 
policy to handle it.
   
   Using an unchecked connector exception preserves the existing public 
`SeaTunnelRowDeserializer` method signature. No configuration option, default 
value, dependency, or valid-message behavior needs to change. The error text 
also avoids including the message body, receipt handle, queue URL, or 
credentials.
   
   ## SeaTunnel Version
   
   Current `dev` at `5dbfb374f985349aeefde9cf84169aa98b3ac5ca`.
   
   ## SeaTunnel Config
   
   ```conf
   env {
     parallelism = 1
     job.mode = "BATCH"
   }
   
   source {
     AmazonSqs {
       url = "http://localhost:4566/000000000000/source_queue";
       access_key_id = "test"
       secret_access_key = "test"
       region = "us-east-1"
       format = json
       delete_message = true
       schema = {
         fields {
           name = string
         }
       }
     }
   }
   
   sink {
     Console {}
   }
   ```
   
   This config shows the affected production setting. The deterministic 
reproduction uses the real `AmazonSqsSourceReader` with an in-memory 
`SqsClient` boundary so the receipt handle passed to `deleteMessage` can be 
asserted without an external service.
   
   ## Running Command
   
   The focused production-reader regression can be run with:
   
   ```shell
   ./mvnw -pl seatunnel-connectors-v2/connector-amazonsqs \
     
-Dtest=AmazonSqsSourceReaderTest#shouldNotDeleteMessageWhenDeserializationThrows
 \
     test
   ```
   
   ## Error Exception
   
   No deserialization exception is propagated by the current source. In the 
pre-fix regression, the expected failure was missing:
   
   ```log
   Expected java.io.IOException to be thrown, but nothing was thrown.
   ```
   
   The same regression recorded the invalid message's receipt handle in the SQS 
deletion call.
   
   ## Zeta or Flink or Spark Version
   
   The defect was reproduced through the connector's production 
`AmazonSqsSourceReader` call path. The Zeta collector path was also inspected 
because it permits the `null` value to advance before deletion.
   
   ## Java or Scala Version
   
   Reproduced independently with:
   
   - Oracle Java 8 (`1.8.0_172`)
   - Eclipse Temurin Java 11 (`11.0.19`)
   
   The completed connector test module passes on both runtimes: 7 tests, 0 
failures, 0 errors on each.
   
   ## Screenshots
   
   Not applicable.
   
   ## 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