
AWS Trigger
CertifiedTrigger on Kinesis records (polling)
AWS Trigger
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.
type: io.kestra.plugin.aws.kinesis.TriggerExamples
Poll a Kinesis stream every 30 seconds
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
streamName *Requiredstring
Stream name
Name of the Kinesis stream to poll.
accessKeyId string
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.
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).
interval Non-dynamicstring
PT1MdurationInterval
iteratorType string
LATESTAT_SEQUENCE_NUMBERAFTER_SEQUENCE_NUMBERTRIM_HORIZONLATESTAT_TIMESTAMPIterator type
Start position: LATEST, TRIM_HORIZON, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER.
maxDuration string
PT30SMax duration
maxRecords integerstring
1000Max records
pollDuration string
PT1SPoll duration
region string
AWS region with which the SDK should communicate
secretKeyId string
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.
sessionToken string
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.
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
count integer
0Record count
Total records consumed.
lastSequencePerShard object
Last sequence per shard
Map of shardId to last consumed sequence number.
uri string
uriRecords file URI
Internal storage URI containing the consumed records (ION).