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]