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

log

Classes

KafkaBrokerPayload

The broker_payload carried alongside each raw event from open_stream.

KafkaSharedStreamProducer

Broker-side half of a shared Kafka stream running in ack mode.

KafkaSharedStreamTrigger

Event-driven trigger that watches Kafka topics through a shared consumer.

Module Contents

airflow.providers.apache.kafka.triggers.shared_stream.log[source]
class airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload[source]

Bases: NamedTuple

The broker_payload carried alongside each raw event from open_stream.

(topic, partition, offset) let KafkaSharedStreamProducer.advance() commit the right TopicPartition and KafkaSharedStreamProducer.get_advance_lane() key the advance by partition. value / key are populated only when a dlq_topic is configured – they are what a dead-letter send re-publishes, so retaining them through the whole outstanding window is pointless without a DLQ.

topic: str[source]
partition: int[source]
offset: int[source]
value: Any = None[source]
key: Any = None[source]
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.SharedStreamProducer

Broker-side half of a shared Kafka stream running in ack mode.

Drives one confluent-kafka Consumer for a shared-stream group and commits offsets only after every subscriber that derived a TriggerEvent from 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_id must set enable.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). Set dlq_topic to 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 poll call.

  • dlq_topic (str | None) – Optional dead-letter topic for rejected messages, produced through kafka_config_id. When unset, rejected messages are dropped.

topics[source]
kafka_config_id = 'kafka_default'[source]
poll_timeout = 1.0[source]
dlq_topic = None[source]
async open_stream()[source]

Open the consumer lazily and yield (value, KafkaBrokerPayload) per message.

get_advance_lane(broker_payload)[source]

Order commits per (topic, partition).

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 offset N + 1 marks everything up to N as 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. If dlq_topic is 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 and seek the 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 offset H, no later batch may commit past H until a re-read batch starting at or before H acks 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.

async aclose()[source]

Flush the DLQ producer and close the consumer when the poll ends; best-effort.

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

Event-driven trigger that watches Kafka topics through a shared consumer.

Triggers that declare the same topics + kafka_config_id share one underlying Kafka consumer in the triggerer (a single poll loop broadcast to every subscriber). Each subscriber fires a TriggerEvent per message; override filter_shared_stream() to fire only for the messages this trigger cares about.

Designed to back an AssetWatcher for event-driven scheduling. The offset is committed only after the derived TriggerEvent is persisted – see KafkaSharedStreamProducer for the enable.auto.commit=false requirement.

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 poll call.

  • dlq_topic (str | None) – Optional dead-letter topic for messages a subscriber rejects; see KafkaSharedStreamProducer. When unset, rejected messages are dropped.

topics[source]
kafka_config_id = 'kafka_default'[source]
poll_timeout = 1.0[source]
dlq_topic = None[source]
serialize()[source]

Return the information needed to reconstruct this Trigger.

Returns:

Tuple of (class path, keyword arguments needed to re-instantiate).

Return type:

tuple[str, dict[str, Any]]

shared_stream_key()[source]

Triggers on the same topics + connection share one consumer.

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 SharedStreamProducer owns the broker connection for the lifetime of one poll — it supplies events through open_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’s open_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 TriggerEvent per message. Override to filter or transform.

abstract run()[source]
Async:

Not supported – this trigger runs only through the shared-stream manager.

shared_stream_key always returns non-None, so the triggerer drives this trigger through filter_shared_stream(); the _SharedStreamGroup owns 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?