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.

Enable compatibility mode

Use it to connect to S3 bucket with S3 compatible services that don't support the new transport client.

The endpoint with which the SDK should communicate

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

Force path style access

Must only be used when compatibilityMode is enabled.

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.

Reference (ref) of the pluginDefaults to apply to this task.

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