goutamadwant opened a new pull request, #12006:
URL: https://github.com/apache/seatunnel/pull/12006
### Purpose of this pull request
Related to #10753.
This PR adds an Azure Queue Storage source to the existing
`connector-azure-queue-storage` module.
The source:
- reads messages from one Azure Storage queue in bounded batches
- supports connection string, shared key, and SAS token authentication
- supports JSON and delimited text payloads
- supports plain and Base64 message encoding
- renews message visibility while messages wait for checkpoint completion
- deletes messages only after the checkpoint containing them completes
- retains messages after aborted checkpoints for a later checkpoint
- bounds retained messages through `max_in_flight_messages`
- releases malformed and unprocessed messages back to the queue
- reports receive, delete, and visibility-renewal failures to the source task
- provides at-least-once delivery semantics
The Azure client construction and authentication validation are shared with
the existing sink. Existing sink behavior is unchanged.
This PR also adds source registration, English and Chinese documentation,
configuration examples, unit tests, and Azurite E2E coverage.
### Does this PR introduce _any_ user-facing change?
Yes.
Users can now configure `AzureQueueStorage` as a streaming source.
The source requires checkpointing because messages are deleted from Azure
Queue Storage only after checkpoint completion. Until then, their visibility
timeout is renewed. A task failure or lease loss can cause redelivery, so the
source provides at-least-once delivery.
One Azure queue is represented by one source split in this first
implementation. Increasing job parallelism does not create additional consumers
for the same queue.
The existing Azure Queue Storage sink configuration and behavior are
unchanged.
### How was this patch tested?
Added unit coverage for:
- source configuration and authentication validation
- streaming and checkpoint requirements
- source factory creation
- checkpoint-complete deletion
- aborted and overlapping checkpoints
- partial delete failure and retry
- visibility renewal and renewed pop-receipt usage
- asynchronous visibility-renewal failure propagation
- bounded in-flight message handling
- malformed payload release
- collector-based deserialization
- reader shutdown and outstanding-message release
- existing sink behavior after sharing client construction
Verified the connector module with JDK 8 and JDK 11:
```shell
./mvnw -pl seatunnel-connectors-v2/connector-azure-queue-storage test
```
Result: 36 tests passed with no failures or errors.
Verified the connector package with JDK 11:
```shell
./mvnw -pl seatunnel-connectors-v2/connector-azure-queue-storage package
```
Verified E2E test compilation:
```shell
./mvnw -pl
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azure-queue-storage-e2e
-DskipTests test-compile
```
The existing Azurite E2E suite now includes a source case that emits a
message, waits for a completed checkpoint, and verifies that the message is
deleted afterward.
### Check list
* [x] No new Jar binary package is added, so no License Notice update is
required.
* [x] English and Chinese source documentation and changelog entries are
added.
* [x] No incompatible change is introduced, so `incompatible-changes.md`
does not require an update.
* [x] Connector integration files were checked:
1. `plugin-mapping.properties` includes the Azure Queue Storage source.
2. `seatunnel-dist` already includes the connector module for the existing
sink.
3. The existing CI label scope already covers the connector module.
4. The existing Azure Queue Storage E2E module now covers the source.
5. `config/plugin_config` already includes the shared connector module.
--
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]