Skip to main content
Version: Next

AzureEventHubs

Azure Event Hubs source connector

Description​

Reads events from one Azure Event Hub through the native Azure AMQP client and converts each event body to a SeaTunnel row.

Use the native connector when the job needs Event Hubs partition discovery and SeaTunnel-managed sequence-number recovery. Azure Event Hubs also exposes a Kafka-compatible endpoint; use the SeaTunnel Kafka connector instead when an existing deployment already standardizes on Kafka protocol configuration and semantics.

Support Those Engines​

Spark
Flink
SeaTunnel Zeta

Key Features​

Options​

nametyperequireddefault value
connection_stringstringyes-
event_hub_namestringyes-
consumer_groupstringno$Default
start_modeenumnoearliest
formatenumnojson
field_delimiterstringno,
max_batch_sizeintno100
poll_timeout_mslongno1000
prefetch_countintno300
schemaconfigyes-
common-optionsno-

connection_string [string]​

Azure Event Hubs namespace connection string. Configure event_hub_name separately; connection strings containing an EntityPath segment are rejected so there is one unambiguous event hub selection path. The option is masked in parsed-config logs, but is not encrypted or decrypted by default. To opt into configuration encryption, include connection_string in env.shade.options and use the configured shade to encrypt its value.

This first version supports namespace connection-string authentication. Microsoft Entra ID, managed identity and custom endpoint authentication are not yet supported.

Use a dedicated SAS policy with only the Listen right; do not use RootManageSharedAccessKey for a source job. A namespace-scoped policy grants access across that namespace, while an Event Hub-scoped policy limits access to that hub. See Azure SAS authorization.

Hub-scoped connection strings containing EntityPath cannot be used unchanged. Manually keep the policy name and key, remove the EntityPath segment, and set event_hub_name to that same hub. The connector does not automatically convert or normalize EntityPath, even when it matches event_hub_name. This only changes how the hub name is supplied; it does not broaden the SAS policy's permissions. The emulator tests do not verify Azure service-side SAS authorization, so validate a hub-scoped policy against the target Azure deployment before use.

event_hub_name [string]​

Name of the Event Hub to consume.

consumer_group [string]​

Consumer group used by the source. Each independently checkpointed job should use a dedicated consumer group.

start_mode [enum]​

Position used only when a job starts without restored source state:

  • earliest: start at each partition's current beginning sequence number.
  • latest: start immediately after each partition's last enqueued sequence number.

The enumerator resolves this mode once into a concrete sequence number per partition. A restored job always uses the sequence number stored in its SeaTunnel checkpoint and does not evaluate start_mode again.

format [enum]​

Event body format:

  • json: reads the body as a JSON object.
  • text: splits the body fields with field_delimiter.

field_delimiter [string]​

Field delimiter used when format = text.

max_batch_size [int]​

Maximum events requested from one partition in one poll. The value must be greater than zero.

poll_timeout_ms [long]​

Maximum time one partition poll waits for events. The value must be between 1 and 5000 milliseconds. A bounded timeout lets source shutdown and split changes interrupt idle polling promptly.

prefetch_count [int]​

Maximum events the Azure SDK prefetches for each partition assigned to a source reader. It must be between 1 and 8000 and at least max_batch_size. These bounds are validated when the source configuration is created, before connecting to Azure. A reader can own multiple partitions, so its total client-side buffer is bounded by this value multiplied by its assigned partition count.

schema [config]​

Schema used to deserialize each event body.

common options​

Source plugin common parameters, please refer to Source Common Options for details.

Partition And Recovery Semantics​

The source is streaming-only. At startup, one SeaTunnel source split is created for each Event Hubs partition and assigned with the regular SeaTunnel split owner calculation. Source parallelism can process different partitions concurrently; parallelism greater than the partition count leaves some readers idle.

This first version discovers partitions only during the initial enumeration. Partitions added after the job starts are not picked up dynamically or when restoring existing source state. Discovering them requires starting without restored source state, which reapplies start_mode to every partition and can replay or skip existing data; plan the restart accordingly.

SeaTunnel checkpoint state is the only recovery authority. The connector does not use Azure Blob Storage checkpointing or EventProcessorClient. A split checkpoint stores the next sequence number to read. Events fetched into the reader queue but not emitted before a checkpoint are replayed after recovery, while emitted events advance the split state. This provides at-least-once delivery when checkpointing is enabled.

The connector can run without checkpointing, but a task or job restart then applies start_mode again because no recovery state exists. Downstream processing should tolerate duplicates when checkpointing is enabled.

If Event Hubs retention removes a checkpointed sequence before restore, the source fails instead of silently resetting to earliest or latest. Invalid JSON or text payloads also fail the source task; the last completed checkpoint determines the replay position.

To recover from a trimmed checkpoint position, stop the failing job and start a new job without restoring the old source state. Choose start_mode explicitly: earliest replays all still-retained events, including events already processed, while latest skips the existing backlog. This choice applies to every partition, not just the affected one. Events already removed by retention cannot be recovered from Event Hubs; reconcile missing data from another source if needed. Restarting with the same checkpoint or changing only start_mode does not reset the saved position.

Deserialization and row-emission errors report the partition ID, event sequence number, a safe failure category (I/O or runtime failure), and a bounded chain of exception class names, including wrapped interruption. They do not retain original exception messages, causes or suppressed exceptions, which may contain private event data. The failed event does not advance the checkpointed position. There is no malformed-event skip or dead-letter option in this version; restoring the same checkpoint can encounter the same invalid event again.

Retry And Failure Behavior​

The Azure SDK applies its built-in AMQP retry policy. Outages that outlast that retry budget surface as source task failures. When checkpointing and job recovery are configured, SeaTunnel resumes each partition from its last completed checkpoint. This connector version does not expose Azure SDK retry or backoff settings. Live Azure authorization and outage recovery have not been verified by the local unit tests or emulator coverage.

Task Example​

env {
parallelism = 2
job.mode = "STREAMING"
checkpoint.interval = 10000
}

source {
AzureEventHubs {
connection_string = "Endpoint=sb://my-namespace.servicebus.windows.net/;SharedAccessKeyName=listen;SharedAccessKey=..."
event_hub_name = "events"
consumer_group = "$Default"
start_mode = earliest
format = json
max_batch_size = 100
poll_timeout_ms = 1000
prefetch_count = 300
schema = {
fields {
event_id = string
event_type = string
}
}
}
}

Changelog​

Change Log
ChangeCommitVersion
[Feature][Connector-V2] Add Azure Event Hubs source connector-Next