
NATS Consume
CertifiedConsume NATS JetStream messages
NATS Consume
Consume NATS JetStream messages
Pulls messages from a JetStream subject with explicit acks and writes them to Kestra internal storage. Requires a stream matching the rendered subject; defaults: deliverPolicy=All, pollDuration=PT2S, batchSize=10. Stops when no messages, maxRecords, or maxDuration is reached.
type: io.kestra.plugin.nats.core.ConsumeExamples
Consume messages from any topic subject matching the kestra.> wildcard, using user password authentication.
id: nats_consume_messages
namespace: company.team
tasks:
- id: consume
type: io.kestra.plugin.nats.core.Consume
url: nats://localhost:4222
username: nats_user
password: "{{ secret('NATS_PASSWORD') }}"
subject: kestra.>
durableId: someDurableId
pollDuration: PT5S
Properties
subject *string
1Subject to consume
Rendered subject or wildcard the JetStream stream is bound to.
url *string
1URL to connect to NATS server
The format is (nats://)server_url: port. You can also provide a connection token like so: nats://token@server_url: port
assets
Assets this task consumes as inputs or produces as outputs, for lineage tracking and the asset graph (Enterprise Edition). A flow declaring this property on a task is rejected in the open-source edition.
io.kestra.core.models.assets.AssetsDeclaration
IGNOREFAILWARNAsset failure behavior
Behavior applied to the task state when a declared asset fails to render, emit, or be persisted (e.g. a lock conflict): FAIL escalates it to FAILED, WARN (default) warns it if it would otherwise succeed, IGNORE leaves the state untouched.
Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.
The assets consumed as inputs.
io.kestra.core.models.assets.AssetIdentifier
1The assets produced as outputs.
io.kestra.plugin.ee.assets.Dataset
1150{}1150io.kestra.plugin.ee.assets.File
1150{}1150io.kestra.plugin.ee.assets.Table
1150{}1150io.kestra.plugin.ee.assets.VM
1150{}1150io.kestra.core.models.assets.External
1150{}1150io.kestra.core.models.assets.Custom
11501Custom asset type
{}1150batchSize integer
10>= 1Batch size
Maximum messages fetched per pull; defaults to 10.
creds string
Credentials files authentification
deliverPolicy string
AllAllLastNewByStartSequenceByStartTimeLastPerSubjectDeliver policy
JetStream deliver policy; defaults to All.
durableId string
Durable consumer name
Optional durable name to resume position between runs.
maxDuration string
Max duration
Optional wall-clock duration after which polling stops.
maxRecords integerstring
Max records
Optional cap on total messages; stops once reached.
password string
Plaintext authentication password
pollDuration string
PT2SPoll duration
Wait time per fetch; defaults to PT2S.
since string
Start time
ISO-8601 date-time rendered and parsed to set the deliver start; ignored if null.
token string
Token authentification
username string
Plaintext authentication username
Outputs
messagesCount integer
Messages consumed
Total acknowledged messages during this run.
uri string
uriOutput file URI
Kestra internal storage URI of the ION file containing consumed messages.