Apache Pulsar Reader

Apache Pulsar Reader

Certified

Read messages from Pulsar topics without subscription

Uses a non-durable reader to fetch messages without creating a subscription. Defaults: deserializer STRING, poll timeout 2s, starts at earliest unless positioned otherwise.

yaml
type: io.kestra.plugin.pulsar.Reader
yaml
id: pulsar_reader
namespace: company.team

tasks:
  - id: reader
    type: io.kestra.plugin.pulsar.Reader
    uri: pulsar://localhost:26650
    topic: test_kestra
    deserializer: JSON
Properties
DefaultSTRING
Possible Values
STRINGJSONBYTES

Value deserializer

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.

Authentication token

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

Maximum read duration

Soft timeout evaluated each second; stops when exceeded.

Maximum records before stop

Soft limit evaluated each second; stops after this many messages if set.

Start from specific message ID

Reads from the message immediately after the provided ID. If neither since nor messageId is set, starts at earliest.

Reference (ref) of the pluginDefaults to apply to this task.

DefaultPT2S

Poll wait duration

Maximum time to wait for a new record when none are immediately available.

Topic schema definition

JSON representation of the topic schema when schema enforcement is enabled.

DefaultNONE
Possible Values
NONEAVROJSON

Topic schema type

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

Rollback duration for start position

Finds the latest message published before the given duration (e.g., PT5M starts 5 minutes in the past).

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.

Number of messages consumed

Formaturi

URI of Kestra storage file with consumed messages