Azure RealtimeTrigger

Azure RealtimeTrigger

Certified

Trigger flows from Azure Event Hubs in real time

Starts an EventProcessorClient that emits one execution per event and checkpoints to Blob Storage. Defaults: consumerGroup=$Default, partitionStartingPosition=EARLIEST. Requires checkpointStoreProperties.connectionString and .containerName. Use Trigger for batch polling.

yaml
type: io.kestra.plugin.azure.eventhubs.RealtimeTrigger

Trigger flow based on events received from Azure Event Hubs in real-time.

yaml
id: azure_eventhubs_realtime_trigger
namespace: company.team

tasks:
  - id: log
    type: io.kestra.plugin.core.log.Log
    message: Hello there! I received {{ trigger.body }} from Azure EventHubs!

triggers:
  - id: read_from_eventhub
    type: io.kestra.plugin.azure.eventhubs.RealtimeTrigger
    eventHubName: my_eventhub
    namespace: my_eventhub_namespace
    connectionString: "{{ secret('EVENTHUBS_CONNECTION') }}"
    bodyDeserializer: JSON
    consumerGroup: "$Default"
    checkpointStoreProperties:
      containerName: kestra
      connectionString: "{{ secret('BLOB_CONNECTION') }}"

Use the Azure Event Hubs Realtime Trigger to push events into Azure Table Storage

yaml
    id: eventhubs_realtime_trigger
    namespace: company.team

    tasks:
      - id: insert_into_storagetable
        type: io.kestra.plugin.azure.storage.table.Bulk
        endpoint: https://yourstorageaccount.table.core.windows.net
        connectionString: "{{ secret('STORAGETABLE_CONNECTION') }}"
        table: orders
        from:
          - partitionKey: order_id
            rowKey: "{{ trigger.body | jq('.order_id') | first }}"
            properties:
              customer_name: "{{ trigger.body | jq('.customer_name') | first }}"
              customer_email: "{{ trigger.body | jq('.customer_email') | first }}"
              product_id: "{{ trigger.body | jq('.product_id') | first }}"
              price: "{{ trigger.body | jq('.price') | first }}"
              quantity: "{{ trigger.body | jq('.quantity') | first }}"
              total: "{{ trigger.body | jq('.total') | first }}"

    triggers:
      - id: realtime_trigger
        type: io.kestra.plugin.azure.eventhubs.RealtimeTrigger
        eventHubName: orders
        namespace: kestra
        connectionString: "{{ secret('EVENTHUBS_CONNECTION') }}"
        bodyDeserializer: JSON
        consumerGroup: $Default
        checkpointStoreProperties:
          containerName: kestra
          connectionString: "{{ secret('BLOB_CONNECTION') }}"
Properties

The event hub to read from

Namespace name of the event hub to connect to

Defaultfalse

Specifies whether a trigger is allowed to start a new execution even if a previous run is still in progress.

DefaultSTRING
Possible Values
STRINGBINARYIONJSON

Body deserializer

Serde used to decode event bodies; defaults to STRING

Default{}

Deserializer properties

Key/value options passed to the selected serde

Default{}

Checkpoint store properties

Blob container config for checkpoints (connectionString, containerName required)

Default5

The maximum number of retry attempts before considering a client operation to have failed

Default500

The maximum permissible delay between retry attempts in milliseconds

Connection string of the Storage Account.

Default$Default

Consumer group

Event Hubs consumer group; defaults to $Default

Custom endpoint address when connecting to the Event Hubs service

Start from enqueue time

ISO-8601 datetime applied only when partitionStartingPosition is set to INSTANT; ignored for EARLIEST and LATEST

DefaultEARLIEST
Possible Values
EARLIESTLATESTINSTANT

Starting position

Initial position strategy per partition; defaults to EARLIEST

The SAS token to use for authenticating requests.

This string should only be the query parameters (with or without a leading '?') and not a full URL.

Shared Key access key for authenticating requests.

Shared Key account name for authenticating requests.

SubTypestring
Possible Values
CREATEDSUBMITTEDRUNNINGPAUSEDRESTARTEDKILLINGSUCCESSWARNINGFAILEDKILLEDCANCELLEDQUEUEDRETRYINGRETRIEDSKIPPEDBREAKPOINTRESUBMITTED

List of execution states after which a trigger should be stopped (a.k.a. disabled).

Defaulttrue

A condition that determines whether the trigger should run.

A Pebble expression evaluated at trigger time. The trigger fires only when the expression evaluates to a truthy value (true, a non-empty string, a non-zero number). Use this to gate trigger execution on dynamic runtime values such as execution labels, flow variables, or environment conditions.

Body

Content Type

Correlation Id

Enqueued Timestamp

Message Id

Offset

Partition Key

Properties

Sequence Number