
Apache Pulsar Trigger
CertifiedTrigger a flow by polling Pulsar messages
Apache Pulsar Trigger
Trigger a flow by polling Pulsar messages
Periodically consumes from topics, stores the batch in Kestra storage, and exposes it via {{ trigger.uri }}. Use RealtimeTrigger for one execution per message.
type: io.kestra.plugin.pulsar.TriggerExamples
id: pulsar_trigger
namespace: company.team
tasks:
- id: log
type: io.kestra.plugin.core.log.Log
message: "{{ trigger.value }}"
triggers:
- id: trigger
type: io.kestra.plugin.pulsar.Trigger
interval: PT30S
topic: kestra_trigger
uri: pulsar://localhost:26650
deserializer: JSON
subscriptionName: kestra_trigger_sub
Properties
deserializer *Requiredstring
STRINGSTRINGJSONBYTESValue deserializer
subscriptionName *Requiredstring
Subscription name
Identifies the subscription so only unconsumed records for that subscription are fetched.
topic *Requiredobject
Source topic(s)
Single topic or list of topics to consume.
uri *Requiredstring
Pulsar service URL
One or more Pulsar protocol URLs, e.g. pulsar://localhost: 6650 or pulsar://host1: 6650,host2: 6651. Use pulsar+ssl:// when enabling TLS.
allowConcurrent Non-dynamicboolean
falseSpecifies whether a trigger is allowed to start a new execution even if a previous run is still in progress.
authenticationToken string
Authentication token
Token used when the broker requires token-based auth (e.g., hosted providers).
consumerName string
Consumer name
Optional name reused on reconnects; helps coordinate with broker policies.
consumerProperties object
Consumer properties
Key/value properties applied to the Pulsar consumer builder.
encryptionKey string
Public encryption key
Key used for payload encryption/decryption when the topic is secured.
initialPosition string
EarliestLatestEarliestInitial subscription position
Where to start consuming (default Earliest).
interval Non-dynamicstring
PT1MdurationInterval between polling.
The interval between 2 different polls of schedule, this can avoid to overload the remote system with too many calls. For most of the triggers that depend on external systems, a minimal interval must be at least PT30S. See ISO_8601 Durations for more information of available interval values.
maxDuration string
Maximum read duration
Soft timeout evaluated each second; stops when exceeded.
maxRecords integerstring
Maximum records before stop
Soft limit evaluated each second; stops after this many messages if set.
pollDuration string
PT2SPoll wait duration
Maximum time to wait for a new record when none are immediately available.
schemaString string
Topic schema definition
JSON schema used when schema enforcement is enabled.
schemaType string
NONENONEAVROJSONTopic schema type
One of NONE (default), AVRO, or JSON.
stopAfter Non-dynamicarray
CREATEDSUBMITTEDRUNNINGPAUSEDRESTARTEDKILLINGSUCCESSWARNINGFAILEDKILLEDCANCELLEDQUEUEDRETRYINGRETRIEDSKIPPEDBREAKPOINTRESUBMITTEDList of execution states after which a trigger should be stopped (a.k.a. disabled).
subscriptionType string
ExclusiveExclusiveSharedFailoverKey_SharedSubscription type
Delivery semantics such as Exclusive (default), Shared, or Failover.
tlsOptions Non-dynamic
TLS options
Certificate/key material for TLS client authentication. Requires a pulsar+ssl:// URL.
io.kestra.plugin.pulsar.AbstractPulsarConnection-TlsOptions
CA certificate
Base64-encoded PEM of the trusted CA chain.
Client certificate
Base64-encoded PEM content for the client certificate.
Client key
Base64-encoded PEM private key matching the client certificate.
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
messagesCount integer
Number of messages consumed
uri string
uriURI of Kestra storage file with consumed messages