Apache Pulsar RealtimeTrigger

Apache Pulsar RealtimeTrigger

Certified

Trigger a flow for each Pulsar message

Consumes messages in real time and emits one execution per message. Use Trigger instead when batching messages over an interval.

yaml
type: io.kestra.plugin.pulsar.RealtimeTrigger

Consume a message from a Pulsar topic in real-time.

yaml
id: pulsar
namespace: company.team

tasks:
  - id: log
    type: io.kestra.plugin.core.log.Log
    message: "{{ trigger.value }}"

triggers:
  - id: realtime_trigger
    type: io.kestra.plugin.pulsar.RealtimeTrigger
    topic: kestra_trigger
    uri: pulsar://localhost:26650
    deserializer: JSON
    subscriptionName: kestra_trigger_sub
Properties
DefaultSTRING
Possible Values
STRINGJSONBYTES

Value deserializer

Subscription name

Identifies the subscription so only unconsumed records for that subscription are fetched.

Source topic(s)

Single topic or list of topics to consume.

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.

Defaultfalse

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

Authentication token

Token used when the broker requires token-based auth (e.g., hosted providers).

Consumer name

Optional name reused on reconnects; helps coordinate with broker policies.

SubTypestring

Consumer properties

Key/value properties applied to the Pulsar consumer builder.

Public encryption key

Key used for payload encryption/decryption when the topic is secured.

DefaultEarliest
Possible Values
LatestEarliest

Initial subscription position

Where to start consuming (default Earliest).

Topic schema definition

JSON schema used when schema enforcement is enabled.

DefaultNONE
Possible Values
NONEAVROJSON

Topic schema type

One of NONE (default), AVRO, or JSON.

SubTypestring
Possible Values
CREATEDSUBMITTEDRUNNINGPAUSEDRESTARTEDKILLINGSUCCESSWARNINGFAILEDKILLEDCANCELLEDQUEUEDRETRYINGRETRIEDSKIPPEDBREAKPOINTRESUBMITTED

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

DefaultExclusive
Possible Values
ExclusiveSharedFailoverKey_Shared

Subscription type

Delivery semantics such as Exclusive (default), Shared, or Failover.

TLS options

Certificate/key material for TLS client authentication. Requires a pulsar+ssl:// URL.

Definitions
castring

CA certificate

Base64-encoded PEM of the trusted CA chain.

certstring

Client certificate

Base64-encoded PEM content for the client certificate.

keystring

Client key

Base64-encoded PEM private key matching the client certificate.

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.

Formatdate-time

The message event time

The message key

The message id

SubTypestring

The message properties

The topic the message belongs to

The message value