AWS RealtimeTrigger

AWS RealtimeTrigger

Certified

Trigger on Kinesis records (realtime)

Subscribes to a Kinesis stream consumer and emits executions as records arrive using SubscribeToShard. Handles shard discovery on a schedule.

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

Realtime Kinesis trigger

yaml
id: realtime_kinesis
namespace: company.team

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

triggers:
  - id: realtime
    type: io.kestra.plugin.aws.kinesis.RealtimeTrigger
    accessKeyId: "{{ secret('AWS_ACCESS_KEY_ID') }}"
    secretKeyId: "{{ secret('AWS_SECRET_KEY_ID') }}"
    region: eu-central-1
    consumerArn: "arn:aws:kinesis:us-east-1:123456789012:stream/my-stream/consumer/kestra-app:1234abcd"
    streamName: "stream"
Properties

Consumer ARN

ARN of a registered stream consumer for enhanced fan-out.

Stream name

Name of the Kinesis stream to subscribe to.

AWS access key ID

Optional static credential. If omitted, the default credentials provider chain is used.

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

DefaultLATEST
Possible Values
AT_SEQUENCE_NUMBERAFTER_SEQUENCE_NUMBERTRIM_HORIZONLATESTAT_TIMESTAMP

Iterator type

Start position: LATEST, TRIM_HORIZON, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER.

AWS region with which the SDK should communicate

AWS secret access key

Pairs with accessKeyId for static credentials. If omitted, the default credentials provider chain is used.

AWS session token for temporary credentials

Used with STS- or SSO-issued temporary credentials. If omitted, the default credentials provider chain is used.

DefaultPT30S

Shard discovery interval

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.

Formatdate-time

Approximate arrival timestamp

The data payload returned by Kinesis

Partition key

Sequence number

Shard ID