AMQP Consume

AMQP Consume

Certified

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.

yaml
type: io.kestra.plugin.amqp.Consume
yaml
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

Broker host

Hostname or IP of the RabbitMQ broker; required unless using the deprecated url.

Queue name to consume

AMQP queue to read from; required and must already exist.

DefaultSTRING
Possible Values
STRINGJSON

Payload serde format

Controls how message bodies are read and written; use STRING for raw text or JSON for structured data. Defaults to STRING.

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.

Definitions
assetFailureBehaviorstring
Possible Values
IGNOREFAILWARN

Asset 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.

enableAutobooleanstring

Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.

inputsarray

The assets consumed as inputs.

id*string
Min length1
typestring
outputs

The assets produced as outputs.

id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*string
Min length1

Custom asset type

descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
Defaultfalse

Automatic 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.

DefaultKestra

Consumer tag

Client-supplied consumer tag used for tracing and cancellations; defaults to Kestra in tasks and triggers.

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.

Maximum records

Soft cap on messages consumed before stopping; evaluated after each ACKed message. Required when maxDuration is not set.

Password

Password for the connection; required when the broker enforces authentication.

Default5672

Broker port

TCP port for AMQP connections; defaults to 5672.

Username

Username for the connection; uses broker default (often guest) when not set.

Default/

Virtual host

Broker virtual host path; defaults to /.

Total messages consumed

Count of messages consumed before the stop condition was reached; acknowledged in a single batch once this output is durably stored.

Formaturi

URI of file storing consumed messages

Internal storage path to the ION-serialized batch returned as taskrun.outputs.uri or trigger.uri.

Unitrecords

The total number of records consumed from the AMQP queue.