AWS Consume

AWS Consume

Certified

Consume records from Kinesis

Reads records from a stream starting at the chosen iterator type. Stops when maxRecords or maxDuration is reached. Writes results to internal storage and records last sequence per shard.

yaml
type: io.kestra.plugin.aws.kinesis.Consume

Consume records from a Kinesis stream using TRIM_HORIZON

yaml
id: kinesis_consume
namespace: company.team

tasks:
  - id: consume
    type: io.kestra.plugin.aws.kinesis.Consume
    accessKeyId: "{{ secret('AWS_ACCESS_KEY_ID') }}"
    secretKeyId: "{{ secret('AWS_SECRET_KEY_ID') }}"
    region: "eu-central-1"
    streamName: "stream"
    iteratorType: TRIM_HORIZON
    pollDuration: PT5S
    maxRecords: 100
Properties

Stream name

Name of the Kinesis stream to read from.

Access Key Id in order to connect to AWS

If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

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

The endpoint with which the SDK should communicate

This property allows you to use a different S3 compatible storage backend.

DefaultLATEST
Possible Values
AT_SEQUENCE_NUMBERAFTER_SEQUENCE_NUMBERTRIM_HORIZONLATESTAT_TIMESTAMP

Iterator type

Start position: LATEST, TRIM_HORIZON, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER, AT_TIMESTAMP.

DefaultPT30S

Max duration

Stop after this duration elapses; default 30s.

Default1000

Max records

Stop after consuming this many records; default 1000.

DefaultPT1S

Poll interval

Sleep between GetRecords calls; default 1s.

AWS region with which the SDK should communicate

Secret Key Id in order to connect to AWS

If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

AWS session token, retrieved from an AWS token service, used for authenticating that this user has received temporary permissions to access a given resource

If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

Starting sequence number

Required when iteratorType is AT_SEQUENCE_NUMBER or AFTER_SEQUENCE_NUMBER.

The AWS STS endpoint with which the SDKClient should communicate

AWS STS Role

The Amazon Resource Name (ARN) of the role to assume. If set the task will use the StsAssumeRoleCredentialsProvider. If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

AWS STS External Id

A unique identifier that might be required when you assume a role in another account. This property is only used when an stsRoleArn is defined.

DefaultPT15M

AWS STS Session duration

The duration of the role session (default: 15 minutes, i.e., PT15M). This property is only used when an stsRoleArn is defined.

AWS STS Session name

This property is only used when an stsRoleArn is defined.

Default0

Record count

Total records consumed.

SubTypestring

Last sequence per shard

Map of shardId to last consumed sequence number.

Formaturi

Records file URI

Internal storage URI containing the consumed records (ION).

The time spent consuming records from the Kinesis stream.

The number of records consumed from the Kinesis stream.