goutamadwant opened a new pull request, #12048:
URL: https://github.com/apache/seatunnel/pull/12048
### Purpose of this pull request
Related to #10753.
This PR adds a native Azure Event Hubs source connector for streaming jobs.
The first source slice includes:
- namespace connection-string authentication with the Event Hub name
configured separately
- one stable SeaTunnel source split per Event Hubs partition
- parallel partition consumption through long-lived, backpressured AMQP
receivers
- `earliest` and `latest` bootstrap modes for fresh jobs
- SeaTunnel checkpoint recovery using the next sequence number for each
partition
- JSON and delimited text deserialization
- bounded batch size, poll timeout, and per-partition prefetch configuration
- connection-string masking in desensitized job configuration
- English and Chinese connector documentation and changelogs
- connector registration, distribution metadata, dependency license entries,
and E2E coverage
#### Recovery and delivery semantics
SeaTunnel checkpoint state is the only recovery authority. The connector
does not use Azure Blob checkpoint storage or `EventProcessorClient`.
`start_mode` is resolved only for a fresh job. Restored readers use the
concrete sequence number stored in the source split checkpoint. The split state
advances only after an event body is successfully deserialized and emitted.
Events fetched but not emitted before a completed checkpoint are replayed
after recovery. This provides at-least-once delivery when checkpointing is
enabled.
The reader validates contiguous sequence numbers before returning records.
If retention removes a checkpointed position, the task fails instead of
silently skipping to the current beginning of the partition.
#### Resource and dependency behavior
Each assigned partition owns one long-lived Azure receiver link.
Demand and the in-memory queue are bounded by `prefetch_count`. The total
reader-side Event Hubs buffer is bounded by `prefetch_count` multiplied by the
number of partitions assigned to that reader.
Azure SDK, Reactor, Reactive Streams, Qpid Proton, and related namespaces
are relocated under the connector-specific shade package.
The packaged connector was inspected to confirm:
- original Azure and AMQP package namespaces do not leak from the connector
JAR
- the SeaTunnel factory service entry is retained
- required third-party dependencies are included in license and
known-dependency metadata
#### First-version scope
- streaming jobs only
- one Event Hub per source instance
- namespace connection-string authentication only
- partitions are discovered during initial enumeration
- partitions added later require a job restart
- no Event Hubs application-property metadata fields
- no Microsoft Entra ID or managed identity authentication
- no custom endpoint
- no timestamp start mode
- no dynamic partition discovery
The Kafka-compatible Event Hubs endpoint remains available through the
existing Kafka connector for deployments using Kafka protocol semantics.
### Does this PR introduce _any_ user-facing change?
Yes.
Users can configure `AzureEventHubs` as a new streaming source.
The connector reads JSON or delimited text event bodies from all partitions
of one Azure Event Hub and restores per-partition sequence positions from
SeaTunnel checkpoints.
This is a new connector and does not change existing connector behavior.
### How was this patch tested?
Added 32 connector unit tests covering:
- default options and positive/negative option validation
- namespace connection-string validation and secret masking
- `earliest` and `latest` initial sequence resolution
- split discovery and deterministic assignment
- pending split checkpoint state
- split return and restore without rediscovery
- restored checkpoint sequence taking priority over `start_mode`
- fair multi-partition polling
- fetch-cursor advancement
- sequence-gap and retention failures without silent skipping
- deserialize-before-checkpoint advancement
- failed-payload replay behavior
- long-lived receiver reuse
- bounded demand and buffering
- buffered-record delivery before terminal errors
- wakeup and close behavior
- streaming-only validation
- factory contract and distributed source serialization
Added 8 focused E2E harness tests covering Azure SDK Reactor thread
lifecycle handling, including overlapping Azure SDK tests and similarly named
non-Azure threads.
Verified connection-string masking through 6 `ConfigBuilderTest` tests.
Verified the connector unit suite on Java 8 and Java 11:
```shell
./mvnw -pl seatunnel-connectors-v2/connector-azure-event-hubs clean test
-DfailIfNoTests=false
```
Verified Java 11 packaging and local installation, including Spotless and
shade processing:
```shell
./mvnw -pl seatunnel-connectors-v2/connector-azure-event-hubs clean install
-DfailIfNoTests=false
```
Verified both affected E2E modules compile after generalizing the Azure SDK
thread lifecycle helper:
```shell
./mvnw -pl
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azure-event-hubs-e2e
-DskipIT=true -DskipUT=true test-compile
./mvnw -pl
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azure-queue-storage-e2e
-DskipIT=true -DskipUT=true test-compile
```
Ran the new E2E against Azurite `3.35.0` and the official Azure Event Hubs
emulator `2.2.1`:
```shell
./mvnw -B -T 1 verify \
-DskipUT=true \
-DskipIT=false \
-Dlicense.skipAddThirdParty=true \
-Dskip.ui=true \
-pl :connector-azure-event-hubs-e2e \
-Pci
```
The E2E passed with two Event Hubs partitions.
It verifies:
- one event is published to each partition
- both rows reach the Console sink
- a SeaTunnel checkpoint completes
- the streaming job can be cancelled
- job and connector resources close cleanly
### Check list
* [x] No new Jar binary is committed. Third-party dependency notices and
known dependency entries are included.
* [x] English and Chinese connector documentation and changelogs are
included.
* [x] No `incompatible-changes.md` update is required because this adds a
new connector without changing an existing contract.
* [x] `plugin-mapping.properties` registers
`seatunnel.source.AzureEventHubs`.
* [x] `seatunnel-dist/pom.xml` includes the connector.
* [x] `.github/workflows/labeler/label-scope-conf.yml` includes the
connector scope.
* [x] An official-emulator E2E module is included under `seatunnel-e2e`.
* [x] `config/plugin_config` includes the connector.
--
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]