AWS Trigger

AWS Trigger

Certified

Trigger on Kinesis records (polling)

Polls a Kinesis stream on a fixed interval and starts a flow when new records are found. Stores last sequence per shard to avoid reprocessing.

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

Poll a Kinesis stream every 30 seconds

yaml
id: kinesis_poll
namespace: company.team

tasks:
  - id: log
    type: io.kestra.plugin.core.log.Log
    message: "Consumed {{ trigger.count }} records."

triggers:
  - id: poll
    type: io.kestra.plugin.aws.kinesis.Trigger
    streamName: "stream"
    iteratorType: "LATEST"
    interval: PT30S
    maxRecords: 100
Properties

Stream name

Name of the Kinesis stream to poll.

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.

Defaultfalse

Specifies whether a trigger is allowed to start a new execution even if a previous run is still in progress.

The endpoint with which the SDK should communicate

This property allows you to target a custom or AWS-compatible Kinesis endpoint (for example, a local test endpoint).

DefaultPT1M
Formatduration

Interval

DefaultLATEST
Possible Values
AT_SEQUENCE_NUMBERAFTER_SEQUENCE_NUMBERTRIM_HORIZONLATESTAT_TIMESTAMP

Iterator type

Start position: LATEST, TRIM_HORIZON, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER.

DefaultPT30S

Max duration

Default1000

Max records

DefaultPT1S

Poll duration

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.

SubTypestring
Possible Values
CREATEDSUBMITTEDRUNNINGPAUSEDRESTARTEDKILLINGSUCCESSWARNINGFAILEDKILLEDCANCELLEDQUEUEDRETRYINGRETRIEDSKIPPEDBREAKPOINTRESUBMITTED

List of execution states after which a trigger should be stopped (a.k.a. disabled).

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.

Defaulttrue

A condition that determines whether the trigger should run.

A Pebble expression evaluated at trigger time. The trigger fires only when the expression evaluates to a truthy value (true, a non-empty string, a non-zero number). Use this to gate trigger execution on dynamic runtime values such as execution labels, flow variables, or environment conditions.

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