Apache Software Foundation / SeaTunnel
Merged upstreamFeatureSelected contributionMerged Sep 28, 2026

Add a native Azure Event Hubs source

Implemented and registered a streaming Azure Event Hubs source with one split per partition, JSON and text bodies, bounded reactive receivers, bilingual documentation, emulator coverage, and recoverable sequence-number state.

apache/seatunnel · #12048

Cloud-streaming feature

Apache SeaTunnel now has a native Azure Event Hubs source with partition-aware parallelism, bounded AMQP backpressure, secure diagnostics, and checkpoint-based at-least-once recovery.

Problem

SeaTunnel lacked a native Event Hubs source that could discover partitions, distribute them across readers, apply backpressure, deserialize records, and restore a precise per-partition position through SeaTunnel checkpoints.

Approach

Adds connector configuration and validation, split enumeration and serialization, partition-scoped AMQP receivers, fetch and checkpoint cursor separation, secure record emission, format-aware deserialization, shaded packaging, distribution registration, documentation, and an official-emulator E2E module.

Impact and scope

  • Maps each Event Hubs partition to a SeaTunnel split, enabling parallel consumption while bounding each receiver by its configured prefetch capacity.
  • Stores the next sequence number only after successful deserialization and emission, so fetched but uncheckpointed events may replay rather than be silently lost after recovery.
  • Fails explicitly on retention gaps and exhausted SDK retries, preserving partition and sequence context without exposing event payloads, credentials, or exception messages that may contain sensitive data.
  • Supports JSON and text payloads, earliest and latest start modes, English and Chinese documentation, and standard SeaTunnel plugin discovery and distribution packaging.
  • Version-one boundaries remain explicit: restored state does not discover new partitions, and the connector does not add Blob checkpoint storage, timestamp starts, Entra ID, managed identity, custom endpoints, or connector-level retry tuning.

Validation

  • On both Java 8 and Java 11, 46 Event Hubs tests and 24 shared configuration tests passed; Java 11 also passed 12 adjacent Azure Queue configuration tests and four documentation checks.
  • Focused verification, formatting, shaded packaging, jar inventory, and isolated ServiceLoader probes validated connector-factory and relocated Jackson-provider loading.
  • The PR includes official-emulator E2E coverage, but the final follow-up did not rerun it; live Azure SAS authorization and outage recovery remain unverified.
  • All four final hosted checks passed, two upstream contributors approved the final authored head, and the GitHub-verified merge credits Goutam Adwant.