
AWS RealtimeTrigger
CertifiedTrigger on Kinesis records (realtime)
AWS RealtimeTrigger
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.
type: io.kestra.plugin.aws.kinesis.RealtimeTriggerExamples
Realtime Kinesis trigger
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
consumerArn *Requiredstring
Consumer ARN
ARN of a registered stream consumer for enhanced fan-out.
streamName *Requiredstring
Stream name
Name of the Kinesis stream to subscribe to.
accessKeyId string
AWS access key ID
Optional static credential. If omitted, the default credentials provider chain is used.
allowConcurrent Non-dynamicboolean
falseSpecifies whether a trigger is allowed to start a new execution even if a previous run is still in progress.
endpointOverride string
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).
iteratorType string
LATESTAT_SEQUENCE_NUMBERAFTER_SEQUENCE_NUMBERTRIM_HORIZONLATESTAT_TIMESTAMPIterator type
Start position: LATEST, TRIM_HORIZON, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER.
region string
AWS region with which the SDK should communicate
secretKeyId string
AWS secret access key
Pairs with accessKeyId for static credentials. If omitted, the default credentials provider chain is used.
sessionToken string
AWS session token for temporary credentials
Used with STS- or SSO-issued temporary credentials. If omitted, the default credentials provider chain is used.
shardDiscoveryInterval string
PT30SShard discovery interval
startingSequenceNumber string
Starting sequence number
Required when iteratorType is AT_SEQUENCE_NUMBER or AFTER_SEQUENCE_NUMBER.
stopAfter Non-dynamicarray
CREATEDSUBMITTEDRUNNINGPAUSEDRESTARTEDKILLINGSUCCESSWARNINGFAILEDKILLEDCANCELLEDQUEUEDRETRYINGRETRIEDSKIPPEDBREAKPOINTRESUBMITTEDList of execution states after which a trigger should be stopped (a.k.a. disabled).
stsEndpointOverride string
The AWS STS endpoint with which the SDKClient should communicate
stsRoleArn string
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.
stsRoleExternalId string
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.
stsRoleSessionDuration string
PT15MAWS 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.
stsRoleSessionName string
AWS STS Session name
This property is only used when an stsRoleArn is defined.
when string
trueA 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.
Outputs
approximateArrivalTimestamp string
date-timeApproximate arrival timestamp
data string
The data payload returned by Kinesis
partitionKey string
Partition key
sequenceNumber string
Sequence number
shardId string
Shard ID