
AMQP Consume
CertifiedConsume AMQP messages until a stop condition
AMQP Consume
Consume AMQP messages until a stop condition
Consumes from a queue, writing messages to internal storage and returning their URI; requires maxDuration or maxRecords to stop. Messages are ACKed in a single batch once the output is durably stored, and NACKed individually on processing failure. Defaults to consumer tag Kestra and serde type STRING.
type: io.kestra.plugin.amqp.ConsumeExamples
id: amqp_consume
namespace: company.team
tasks:
- id: consume
type: io.kestra.plugin.amqp.Consume
host: localhost
port: 5672
username: guest
password: "{{ secret('AMQP_PASSWORD') }}"
virtualHost: /my_vhost
queue: kestramqp.queue
maxRecords: 1000
Properties
host *string
Broker host
Hostname or IP of the RabbitMQ broker; required unless using the deprecated url.
queue *string
Queue name to consume
AMQP queue to read from; required and must already exist.
serdeType *string
STRINGSTRINGJSONPayload serde format
Controls how message bodies are read and written; use STRING for raw text or JSON for structured data. Defaults to STRING.
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
{}1150autoAck booleanstring
falseAutomatic acknowledgment
When true, the broker acknowledges messages as soon as they are delivered; they are not requeued if the task is killed mid-batch. When false, the task sends a single bulk acknowledgment for the whole batch once it is durably stored, and NACKs a message on processing failure; if killed before that point, the broker requeues every unacknowledged message.
consumerTag string
KestraConsumer tag
Client-supplied consumer tag used for tracing and cancellations; defaults to Kestra in tasks and triggers.
maxDuration string
Maximum duration
Soft cap on run time; checked roughly every 100 ms so actual runtime can slightly exceed this value. Required when maxRecords is not set.
maxRecords integerstring
Maximum records
Soft cap on messages consumed before stopping; evaluated after each ACKed message. Required when maxDuration is not set.
password string
Password
Password for the connection; required when the broker enforces authentication.
port string
5672Broker port
TCP port for AMQP connections; defaults to 5672.
username string
Username
Username for the connection; uses broker default (often guest) when not set.
virtualHost string
/Virtual host
Broker virtual host path; defaults to /.
Outputs
count integer
Total messages consumed
Count of messages consumed before the stop condition was reached; acknowledged in a single batch once this output is durably stored.
uri string
uriURI of file storing consumed messages
Internal storage path to the ION-serialized batch returned as taskrun.outputs.uri or trigger.uri.
Metrics
consumed.records counter
recordsThe total number of records consumed from the AMQP queue.