
Azure Trigger
CertifiedPoll Azure Event Hubs and trigger flows
Azure Trigger
Poll Azure Event Hubs and trigger flows
Periodically consumes events in batches, checkpoints to Blob Storage, and triggers one execution per batch. Defaults: interval=PT60S, consumerGroup=$Default, partitionStartingPosition=EARLIEST, maxBatchSizePerPartition=50, maxWaitTimePerPartition=PT5S, maxDuration=PT10S. Use RealtimeTrigger for per-event executions.
type: io.kestra.plugin.azure.eventhubs.TriggerExamples
Trigger flow based on events received from Azure Event Hubs in batch.
id: azure_eventhubs_trigger
namespace: company.team
tasks:
- id: log
type: io.kestra.plugin.core.log.Log
message: Hello there! I received {{ trigger.eventsCount }} from Azure EventHubs!
triggers:
- id: read_from_eventhub
type: io.kestra.plugin.azure.eventhubs.Trigger
interval: PT30S
eventHubName: my_eventhub
namespace: my_eventhub_namespace
connectionString: "{{ secret('EVENTHUBS_CONNECTION') }}"
bodyDeserializer: JSON
consumerGroup: "$Default"
checkpointStoreProperties:
containerName: kestra
connectionString: "{{ secret('BLOB_CONNECTION') }}"
Properties
checkpointStoreProperties *Requiredobject
Checkpoint store properties
Blob container config for checkpoints (connectionString, containerName required)
eventHubName *Requiredstring
The event hub to read from
namespace *Requiredstring
Namespace name of the event hub to connect to
allowConcurrent Non-dynamicboolean
falseSpecifies whether a trigger is allowed to start a new execution even if a previous run is still in progress.
bodyDeserializer string
STRINGSTRINGBINARYIONJSONBody deserializer
Serde used to decode event bodies; defaults to STRING
bodyDeserializerProperties object
{}Deserializer properties
Key/value options passed to the selected serde
clientMaxRetries integerstring
5The maximum number of retry attempts before considering a client operation to have failed
clientRetryDelay integerstring
500The maximum permissible delay between retry attempts in milliseconds
connectionString string
Connection string of the Storage Account.
consumerGroup string
$DefaultConsumer group
Event Hubs consumer group; defaults to $Default
customEndpointAddress string
Custom endpoint address when connecting to the Event Hubs service
enqueueTime string
Start from enqueue time
ISO-8601 datetime applied only when partitionStartingPosition is set to INSTANT; ignored for EARLIEST and LATEST
interval Non-dynamicstring
PT1MdurationPolling interval
Time between poll cycles; defaults to PT60S
maxBatchSizePerPartition integerstring
50Max batch size per partition
Maximum events pulled per partition read; defaults to 50
maxDuration string
PT10SOverall max duration
Stop consuming after this duration each poll; defaults to PT10S
maxWaitTimePerPartition string
PT5SMax wait per partition
Maximum wait for a partition batch before returning; defaults to PT5S
partitionStartingPosition string
EARLIESTEARLIESTLATESTINSTANTStarting position
Initial position strategy per partition; defaults to EARLIEST
sasToken string
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.
stopAfter Non-dynamicarray
CREATEDSUBMITTEDRUNNINGPAUSEDRESTARTEDKILLINGSUCCESSWARNINGFAILEDKILLEDCANCELLEDQUEUEDRETRYINGRETRIEDSKIPPEDBREAKPOINTRESUBMITTEDList of execution states after which a trigger should be stopped (a.k.a. disabled).
when string
trueA 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.
Outputs
eventsCount integer
Events consumed
uri string
uriConsumed events URI
kestra:// URI for the ION file containing consumed events