RealtimeTrigger
yaml
type: "io.kestra.plugin.kafka.realtimetrigger"
Examples
yaml
id: kafka_realtime_trigger
namespace: company.team
tasks:
- id: log
type: io.kestra.plugin.core.log.Log
message: "{{ trigger.value }}"
triggers:
- id: realtime_trigger
type: io.kestra.plugin.kafka.RealtimeTrigger
topic: test_kestra
properties:
bootstrap.servers: localhost:9092
serdeProperties:
schema.registry.url: http://localhost:8085
keyDeserializer: STRING
valueDeserializer: AVRO
groupId: kafkaConsumerGroupId
yaml
id: kafka_realtime_trigger
namespace: company.team
tasks:
- id: insert_into_mongodb
type: io.kestra.plugin.mongodb.InsertOne
connection:
uri: mongodb://mongoadmin:secret@localhost:27017/?authSource=admin
database: kestra
collection: products
document: |
{
"product_id": "{{ trigger.value | jq('.product_id') | first }}",
"product_name": "{{ trigger.value | jq('.product_name') | first }}",
"category": "{{ trigger.value | jq('.product_category') | first }}",
"brand": "{{ trigger.value | jq('.brand') | first }}"
}
triggers:
- id: realtime_trigger
type: io.kestra.plugin.kafka.RealtimeTrigger
topic: products
properties:
bootstrap.servers: localhost:9092
serdeProperties:
valueDeserializer: JSON
groupId: kestraConsumer
Properties
groupId *Requiredstring
properties *Requiredobject
SubType string
conditions Non-dynamicDateTimeBetweenDayWeekDayWeekInMonthExecutionFlowExecutionLabelsExecutionNamespaceExecutionOutputsExecutionStatusExpressionFlowConditionFlowNamespaceConditionHasRetryAttemptMultipleConditionNotOrPublicHolidayTimeBetweenWeekend
keyDeserializer string
Default
STRING
Possible Values
STRING
INTEGER
FLOAT
DOUBLE
LONG
SHORT
BYTE_ARRAY
BYTE_BUFFER
BYTES
UUID
VOID
AVRO
JSON
onSerdeError Non-dynamicKafkaConsumerInterface-OnSerdeError
partitions array
SubType integer
serdeProperties object
SubType string
Default
{}
since string
stopAfter Non-dynamicarray
SubType string
Possible Values
CREATED
RUNNING
PAUSED
RESTARTED
KILLING
SUCCESS
WARNING
FAILED
KILLED
CANCELLED
QUEUED
RETRYING
RETRIED
SKIPPED
BREAKPOINT
topic Non-dynamicobject
topicPattern string
valueDeserializer string
Default
STRING
Possible Values
STRING
INTEGER
FLOAT
DOUBLE
LONG
SHORT
BYTE_ARRAY
BYTE_BUFFER
BYTES
UUID
VOID
AVRO
JSON
Outputs
Definitions
io.kestra.core.models.triggers.TimeWindow
deadline string
Format
partial-time
endTime string
Format
partial-time
startTime string
Format
partial-time
type string
Default
DURATION_WINDOW
Possible Values
DAILY_TIME_DEADLINE
DAILY_TIME_WINDOW
DURATION_WINDOW
SLIDING_WINDOW
window string
Format
duration
windowAdvance string
Format
duration
Condition for a specific flow of an execution.
flowId *Requiredstring
namespace *Requiredstring
type *Requiredobject
Condition for a flow namespace.
namespace *Requiredstring
type *Requiredobject
prefix boolean
Default
false
Condition for a specific flow. Note that this condition is deprecated, use `io.kestra.plugin.core.condition.ExecutionFlow` instead.
flowId *Requiredstring
namespace *Requiredstring
type *Requiredobject
Condition to allow events between two specific times.
type *Requiredobject
after string
Format
time
before string
Format
time
date string
Default
{{ trigger.date }}
org.apache.commons.lang3.tuple.Pair
Condition that checks labels of an execution.
labels *Requiredarrayobject
type *Requiredobject
Condition based on the outputs of an upstream execution.
expression *Requiredbooleanstring
type *Requiredobject
Condition to allow events on weekend.
type *Requiredobject
date string
Default
{{ trigger.date }}
io.kestra.plugin.kafka.KafkaConsumerInterface-OnSerdeError
topic string
type string
Default
SKIPPED
Possible Values
SKIPPED
DLQ
STORE
Condition to have at least one condition validated.
Condition for an execution namespace.
namespace *Requiredstring
type *Requiredobject
comparison string
Possible Values
EQUALS
PREFIX
SUFFIX
prefix booleanstring
Default
false
Run a flow if the list of preconditions is met in a time window.
id *Requiredstring
Validation RegExp
^[a-zA-Z0-9][a-zA-Z0-9_-]*
Min length
1
type *Requiredobject
resetOnSuccess boolean
Default
true
timeWindow TimeWindow
Default
{
"type": "DURATION_WINDOW"
}
Condition to exclude other conditions.
Condition to execute tasks on a specific day of the week relative to the current month (first, last, ...)
dayInMonth *Requiredstring
Possible Values
FIRST
LAST
SECOND
THIRD
FOURTH
dayOfWeek *Requiredstring
Possible Values
MONDAY
TUESDAY
WEDNESDAY
THURSDAY
FRIDAY
SATURDAY
SUNDAY
type *Requiredobject
date string
Default
{{ trigger.date }}
Condition based on variable expression.
expression *Requiredstring
type *Requiredobject
Condition to allow events on a particular day of the week.
dayOfWeek *Requiredstring
Possible Values
MONDAY
TUESDAY
WEDNESDAY
THURSDAY
FRIDAY
SATURDAY
SUNDAY
type *Requiredobject
date string
Default
{{ trigger.date }}
Condition based on execution status.
type *Requiredobject
in array
SubType string
Possible Values
CREATED
RUNNING
PAUSED
RESTARTED
KILLING
SUCCESS
WARNING
FAILED
KILLED
CANCELLED
QUEUED
RETRYING
RETRIED
SKIPPED
BREAKPOINT
notIn array
SubType string
Possible Values
CREATED
RUNNING
PAUSED
RESTARTED
KILLING
SUCCESS
WARNING
FAILED
KILLED
CANCELLED
QUEUED
RETRYING
RETRIED
SKIPPED
BREAKPOINT
Condition to allow events between two specific datetime values.
type *Requiredobject
after string
Format
date-time
before string
Format
date-time
date string
Default
{{ trigger.date }}
Condition that matches if any taskRun has retry attempts.
type *Requiredobject
in array
SubType string
Possible Values
CREATED
RUNNING
PAUSED
RESTARTED
KILLING
SUCCESS
WARNING
FAILED
KILLED
CANCELLED
QUEUED
RETRYING
RETRIED
SKIPPED
BREAKPOINT
notIn array
SubType string
Possible Values
CREATED
RUNNING
PAUSED
RESTARTED
KILLING
SUCCESS
WARNING
FAILED
KILLED
CANCELLED
QUEUED
RETRYING
RETRIED
SKIPPED
BREAKPOINT
Condition to allow events on public holidays.
type *Requiredobject
country string
date string
Default
{{ trigger.date }}