airflow.providers.apache.kafka.triggers.shared_stream¶
Shared-stream Kafka trigger and producer for event-driven scheduling.
Triggers that declare the same topics + kafka_config_id share a
single Kafka consumer in the triggerer (one poll loop broadcast to every
subscriber) instead of opening one consumer each.
The KafkaSharedStreamProducer owns that consumer and commits
offsets only after the derived TriggerEvent
instances have been persisted, via the shared-stream ack channel.
This module builds on the shared-stream ack channel, which was added in
Airflow 3.3. Importing this module on an older version raises
AirflowOptionalProviderFeatureException.
Attributes¶
Classes¶
The |
|
Broker-side half of a shared Kafka stream running in ack mode. |
|
Event-driven trigger that watches Kafka topics through a shared consumer. |
Module Contents¶
- class airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload[source]¶
Bases:
NamedTupleThe
broker_payloadcarried alongside each raw event fromopen_stream.(topic, partition, offset)letKafkaSharedStreamProducer.advance()commit the rightTopicPartitionandKafkaSharedStreamProducer.get_advance_lane()key the advance by partition.value/keyare populated only when adlq_topicis configured – they are what a dead-letter send re-publishes, so retaining them through the whole outstanding window is pointless without a DLQ.
- class airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer(*, topics, kafka_config_id='kafka_default', poll_timeout=1.0, dlq_topic=None)[source]¶
Bases:
airflow.triggers.shared_stream.SharedStreamProducerBroker-side half of a shared Kafka stream running in ack mode.
Drives one confluent-kafka
Consumerfor a shared-stream group and commits offsets only after every subscriber that derived aTriggerEventfrom a message has had it persisted – the ack channel gates the commit, so a triggerer crash cannot drop a message the broker already considers delivered.Warning
The Kafka connection used by
kafka_config_idmust setenable.auto.commit=false. With auto-commit on, the consumer commits offsets on its own schedule regardless of the ack channel, which defeats the persistence gate and loses messages on a triggerer crash.Events that a subscriber rejects (via
reject_shared_stream_event) are terminal: by default they are dropped (committed past). Setdlq_topicto instead re-publish each rejected message to a dead-letter topic before committing past it. Involuntary failures (ack timeout / overflow) and broadcasts no subscriber was online for are never dropped – the offset floor is held and the consumer is sought back so the message is redelivered this session: a failure to the remaining healthy subscribers, a zero-subscriber broadcast to whoever subscribes next.- Parameters:
topics (collections.abc.Sequence[str]) – Topics the shared consumer subscribes to.
kafka_config_id (str) – Kafka connection id, defaults to
kafka_default.poll_timeout (float) – Seconds the consumer waits on each
pollcall.dlq_topic (str | None) – Optional dead-letter topic for rejected messages, produced through
kafka_config_id. When unset, rejected messages are dropped.
- async open_stream()[source]¶
Open the consumer lazily and yield (value, KafkaBrokerPayload) per message.
- async advance(batch)[source]¶
Commit one partition’s resolved prefix; hold and seek back the first held offset.
Every item shares one
(topic, partition)(the lane). A Kafka commit is cumulative, so committing offsetN + 1marks everything up toNas consumed. Within a batch we commit only through the last item that was terminally handled and stop at the first that should come back:acked– accepted; safe to commit past.rejected– terminally refused. Ifdlq_topicis set the message is re-published there (and flushed) before being committed past; otherwise it is dropped. Either way it is not redelivered.failed(ack timeout / overflow) or all-zero (a broadcast no subscriber was online for) – hold the floor here andseekthe consumer back to it, so it and everything after are redelivered this session. A failure lands on the healthy subscribers the manager left online; a zero-subscriber broadcast reaches whoever subscribes next, instead of waiting for a full consumer rebuild. The floor lifts once a re-read batch acks through the held offset.
The floor is per-lane and carried across calls (
self._floor): once a lane holds at offsetH, no later batch may commit pastHuntil a re-read batch starting at or beforeHacks through it.The batch’s rejected messages are produced to the DLQ together and flushed once, before the offset is committed, so a crash cannot commit past a message that never reached the DLQ. If the commit later fails the batch is redelivered, which may re-send an already-dead-lettered message – DLQ consumers should tolerate duplicates.
- class airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger(*, topics, kafka_config_id='kafka_default', poll_timeout=1.0, dlq_topic=None)[source]¶
Bases:
airflow.triggers.base.BaseEventTriggerEvent-driven trigger that watches Kafka topics through a shared consumer.
Triggers that declare the same
topics+kafka_config_idshare one underlying Kafka consumer in the triggerer (a single poll loop broadcast to every subscriber). Each subscriber fires aTriggerEventper message; overridefilter_shared_stream()to fire only for the messages this trigger cares about.Designed to back an
AssetWatcherfor event-driven scheduling. The offset is committed only after the derivedTriggerEventis persisted – seeKafkaSharedStreamProducerfor theenable.auto.commit=falserequirement.- Parameters:
topics (collections.abc.Sequence[str]) – Topics to watch.
kafka_config_id (str) – Kafka connection id, defaults to
kafka_default.poll_timeout (float) – Seconds the consumer waits on each
pollcall.dlq_topic (str | None) – Optional dead-letter topic for messages a subscriber rejects; see
KafkaSharedStreamProducer. When unset, rejected messages are dropped.
- classmethod create_shared_stream_producer(kwargs)[source]¶
Build the broker-side producer for this trigger’s shared stream (ack mode).
Overriding this classmethod opts the shared stream into ack mode. The manager calls it once per shared-stream group; the returned
SharedStreamProducerowns the broker connection for the lifetime of one poll — it supplies events throughopen_stream, is told to advance the broker in per-lane batches of events whose subscribers have all moved past them and had their derived trigger events confirmed persisted, and is closed when the poll ends. Do not open the broker connection here; open it lazily inside the producer’sopen_stream.Triggers that do not override this method run the fast path: subscribers receive raw events from
open_shared_stream()exactly as before. The stream shape is the same in both modes.
- async filter_shared_stream(shared_stream)[source]¶
Fire one
TriggerEventper message. Override to filter or transform.
- abstract run()[source]¶
- Async:
Not supported – this trigger runs only through the shared-stream manager.
shared_stream_keyalways returns non-None, so the triggerer drives this trigger throughfilter_shared_stream(); the_SharedStreamGroupowns the Kafka consumer and offset commits. There is no standalone path: committing offsets safely needs the ack channel to gate them on trigger-event persistence, which only the manager provides.
Was this entry helpful?