
Apache Pulsar Consume
CertifiedConsume messages from Pulsar topics
Apache Pulsar Consume
Consume messages from Pulsar topics
Subscribes to one or more topics, acknowledges messages after batch receive, and writes them to Kestra storage. Defaults: deserializer STRING, poll timeout 2s, subscription type Exclusive starting at Earliest.
type: io.kestra.plugin.pulsar.ConsumeExamples
id: pulsar_consume
namespace: company.team
tasks:
- id: consume
type: io.kestra.plugin.pulsar.Consume
uri: pulsar://localhost:26650
topic: test_kestra
deserializer: JSON
subscriptionName: kestra_flow
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.
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).
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.
pluginDefaultsRef Non-dynamicstring
Reference (ref) of the pluginDefaults to apply to this task.
pollDuration string
PT2SPoll wait duration
Maximum time to wait for a new record when none are immediately available.
schemaString string
Topic schema definition
JSON representation of the topic schema when schema enforcement is enabled.
schemaType string
NONENONEAVROJSONTopic schema type
One of NONE (default, no enforcement), AVRO, or JSON.
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.
Outputs
messagesCount integer
Number of messages consumed
uri string
uriURI of Kestra storage file with consumed messages