Apache Kafka connection¶
Produce / consume Kafka from a task via a managed Leoflow Connection and the
Apache Kafka provider hooks. The conn_type is kafka (from
apache-airflow-providers-apache-kafka).
Declare the provider¶
Fields the UI asks for¶
A Kafka connection is all Extra: the entire confluent-kafka client config goes into the Extra JSON. There is no host/login โ the brokers and security live in the config dict.
| Field | Value |
|---|---|
| Conn Id | e.g. kafka_default. |
| Extra | The client config, e.g. {"bootstrap.servers": "broker1:9092,broker2:9092", "security.protocol": "SASL_SSL", "sasl.mechanism": "PLAIN", "sasl.username": "etl", "sasl.password": "..."} |
The whole config (including sasl.password) is encrypted at rest (ADR 0019); the
Extra round-trip is pinned by TestKafkaConnectionURIShapeIntegration.
Example DAG (copy-paste)¶
Docs-only recipe (needs a Kafka cluster). The hook is imported inside the task.
# dag.py
from __future__ import annotations
from airflow.sdk import DAG, task
@task
def produce() -> None:
from airflow.providers.apache.kafka.hooks.produce import KafkaProducerHook
hook = KafkaProducerHook(kafka_config_id="kafka_default")
producer = hook.get_producer()
print("produce: sending message to topic 'events'")
producer.produce("events", value=b"hello from leoflow")
producer.flush()
print("produce: ok")
with DAG("kafka_produce", schedule=None, catchup=False, tags=["example"]):
produce()
# leoflow.yaml
schema_version: "1.0"
dag_id: kafka_produce
description: Produce a message to Kafka via KafkaProducerHook.
owner: examples
tags:
- example
python_version: "3.11"
connectors:
- kafka
Run it¶
- Admin โ Connections โ +, type
kafka. Put the client config in Extra. leoflow lite path/to/this/dagโ triggerkafka_produce.
Security notes¶
- The SASL password sits inside the Extra config โ it is encrypted at rest. Prefer
SASL_SSLover plaintext. - Never
print()the Extra โ it carries the SASL credentials.
Related¶
docs/connections/index.mdโ Installing a connector's provider.- ADR 0019 / 0021 โ secret encryption + agent delivery.