Community
Meetups
Documentation
Registry
Use Cases
Announcements
Blog
Ecosystem
Community
Meetups
Documentation
Registry
Use Cases
Announcements
Blog
Ecosystem
Content
Version:
2.1.0
No matching versions
Search docs
⌘
K
Basics
Home
Changelog
Security
Guides
Connection
Hooks
Operators
Message Queues
Sensors
Triggers
References
Configuration
Python API
System tests
System Tests
Resources
Example Dags
PyPI Repository
Installing from sources
Commits
Detailed list of commits
Version:
2.1.0
No matching versions
Search docs
⌘
K
Basics
Home
Changelog
Security
Guides
Connection
Hooks
Operators
Message Queues
Sensors
Triggers
References
Configuration
Python API
System tests
System Tests
Resources
Example Dags
PyPI Repository
Installing from sources
Commits
Detailed list of commits
Home
Index
Index
_
|
A
|
B
|
C
|
D
|
E
|
F
|
G
|
H
|
K
|
L
|
M
|
N
|
O
|
P
|
Q
|
R
|
S
|
T
|
U
|
V
|
X
_
__version__ (in module airflow.providers.apache.kafka)
A
acked() (in module airflow.providers.apache.kafka.operators.produce)
aclose() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer method)
advance() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer method)
airflow.providers.apache.kafka
module
airflow.providers.apache.kafka.assets
module
airflow.providers.apache.kafka.assets.kafka
module
airflow.providers.apache.kafka.get_provider_info
module
airflow.providers.apache.kafka.hooks
module
airflow.providers.apache.kafka.hooks.base
module
airflow.providers.apache.kafka.hooks.client
module
airflow.providers.apache.kafka.hooks.consume
module
airflow.providers.apache.kafka.hooks.produce
module
airflow.providers.apache.kafka.operators
module
airflow.providers.apache.kafka.operators.consume
module
airflow.providers.apache.kafka.operators.produce
module
airflow.providers.apache.kafka.plugins
module
airflow.providers.apache.kafka.plugins.event_producer
module
airflow.providers.apache.kafka.queues
module
airflow.providers.apache.kafka.queues.kafka
module
airflow.providers.apache.kafka.sensors
module
airflow.providers.apache.kafka.sensors.kafka
module
airflow.providers.apache.kafka.triggers
module
airflow.providers.apache.kafka.triggers.await_message
module
airflow.providers.apache.kafka.triggers.msg_queue
module
airflow.providers.apache.kafka.triggers.shared_stream
module
airflow.providers.apache.kafka.version_compat
module
AIRFLOW_V_3_0_PLUS (in module airflow.providers.apache.kafka.version_compat)
AIRFLOW_V_3_1_PLUS (in module airflow.providers.apache.kafka.version_compat)
AIRFLOW_V_3_3_PLUS (in module airflow.providers.apache.kafka.version_compat)
apply_function (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
apply_function() (in module tests.system.apache.kafka.example_dag_kafka_message_queue_trigger)
(in module tests.system.apache.kafka.example_dag_message_queue_trigger)
apply_function_args (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
apply_function_batch (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
apply_function_kwargs (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
asset (in module tests.system.apache.kafka.example_dag_kafka_message_queue_trigger)
(in module tests.system.apache.kafka.example_dag_message_queue_trigger)
await_function() (in module tests.system.apache.kafka.example_dag_event_listener)
(in module tests.system.apache.kafka.example_dag_hello_kafka)
AwaitMessageSensor (class in airflow.providers.apache.kafka.sensors.kafka)
AwaitMessageTrigger (class in airflow.providers.apache.kafka.triggers.await_message)
AwaitMessageTriggerFunctionSensor (class in airflow.providers.apache.kafka.sensors.kafka)
B
BLUE (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
C
CALLBACK_ALLOWLIST_CONFIG_OPTION (in module airflow.providers.apache.kafka.hooks.base)
CALLBACK_CONFIG_KEYS (in module airflow.providers.apache.kafka.hooks.base)
cleanup() (airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger method)
commit_cadence (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
commit_offset (airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
CONFIG_SECTION (in module airflow.providers.apache.kafka.plugins.event_producer)
conn_name_attr (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook attribute)
conn_type (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook attribute)
ConsumeFromTopicOperator (class in airflow.providers.apache.kafka.operators.consume)
consumer_function() (in module tests.system.apache.kafka.example_dag_hello_kafka)
consumer_function_batch() (in module tests.system.apache.kafka.example_dag_hello_kafka)
consumer_logger (in module tests.system.apache.kafka.example_dag_hello_kafka)
convert_asset_to_openlineage() (in module airflow.providers.apache.kafka.assets.kafka)
create_asset() (in module airflow.providers.apache.kafka.assets.kafka)
create_shared_stream_producer() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger class method)
create_topic() (airflow.providers.apache.kafka.hooks.client.KafkaAdminClientHook method)
D
DagRunListener (class in airflow.providers.apache.kafka.plugins.event_producer)
default_args (in module tests.system.apache.kafka.example_dag_hello_kafka)
default_conn_name (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook attribute)
delete_topic() (airflow.providers.apache.kafka.hooks.client.KafkaAdminClientHook method)
delivery_callback (airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
dlq_topic (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger attribute)
E
error_callback() (in module airflow.providers.apache.kafka.hooks.consume)
event_triggered_function (airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
execute() (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator method)
(airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator method)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor method)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor method)
execute_complete() (airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor method)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor method)
F
filter_shared_stream() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger method)
G
get_advance_lane() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer method)
get_conn (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook property)
get_consumer() (airflow.providers.apache.kafka.hooks.consume.KafkaConsumerHook method)
get_producer() (airflow.providers.apache.kafka.hooks.produce.KafkaProducerHook method)
get_provider_info() (in module airflow.providers.apache.kafka.get_provider_info)
get_ui_field_behaviour() (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook class method)
H
hello_kafka() (in module tests.system.apache.kafka.example_dag_hello_kafka)
hook (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator property)
hook_name (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook attribute)
K
KAFKA_COMMON_CONFIG_SECTION (in module airflow.providers.apache.kafka.hooks.base)
kafka_config_id (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook attribute)
(airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger attribute)
KafkaAdminClientHook (class in airflow.providers.apache.kafka.hooks.client)
KafkaAuthenticationError
KafkaBaseHook (class in airflow.providers.apache.kafka.hooks.base)
KafkaBrokerPayload (class in airflow.providers.apache.kafka.triggers.shared_stream)
KafkaConsumerHook (class in airflow.providers.apache.kafka.hooks.consume)
KafkaEventProducerPlugin (class in airflow.providers.apache.kafka.plugins.event_producer)
KafkaMessageQueueProvider (class in airflow.providers.apache.kafka.queues.kafka)
KafkaMessageQueueTrigger (class in airflow.providers.apache.kafka.triggers.msg_queue)
KafkaProducerHook (class in airflow.providers.apache.kafka.hooks.produce)
KafkaSharedStreamProducer (class in airflow.providers.apache.kafka.triggers.shared_stream)
KafkaSharedStreamTrigger (class in airflow.providers.apache.kafka.triggers.shared_stream)
key (airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload attribute)
L
listeners (airflow.providers.apache.kafka.plugins.event_producer.KafkaEventProducerPlugin attribute)
load_connections() (in module tests.system.apache.kafka.example_dag_event_listener)
(in module tests.system.apache.kafka.example_dag_hello_kafka)
local_logger (in module airflow.providers.apache.kafka.operators.produce)
log (in module airflow.providers.apache.kafka.plugins.event_producer)
(in module airflow.providers.apache.kafka.triggers.await_message)
(in module airflow.providers.apache.kafka.triggers.shared_stream)
M
max_batch_size (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
max_messages (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
module
airflow.providers.apache.kafka
airflow.providers.apache.kafka.assets
airflow.providers.apache.kafka.assets.kafka
airflow.providers.apache.kafka.get_provider_info
airflow.providers.apache.kafka.hooks
airflow.providers.apache.kafka.hooks.base
airflow.providers.apache.kafka.hooks.client
airflow.providers.apache.kafka.hooks.consume
airflow.providers.apache.kafka.hooks.produce
airflow.providers.apache.kafka.operators
airflow.providers.apache.kafka.operators.consume
airflow.providers.apache.kafka.operators.produce
airflow.providers.apache.kafka.plugins
airflow.providers.apache.kafka.plugins.event_producer
airflow.providers.apache.kafka.queues
airflow.providers.apache.kafka.queues.kafka
airflow.providers.apache.kafka.sensors
airflow.providers.apache.kafka.sensors.kafka
airflow.providers.apache.kafka.triggers
airflow.providers.apache.kafka.triggers.await_message
airflow.providers.apache.kafka.triggers.msg_queue
airflow.providers.apache.kafka.triggers.shared_stream
airflow.providers.apache.kafka.version_compat
tests.system.apache.kafka
tests.system.apache.kafka.example_dag_event_listener
tests.system.apache.kafka.example_dag_hello_kafka
tests.system.apache.kafka.example_dag_kafka_message_queue_trigger
tests.system.apache.kafka.example_dag_message_queue_trigger
MSK_BOOTSTRAP_SERVERS_REGEX (in module airflow.providers.apache.kafka.hooks.base)
N
name (airflow.providers.apache.kafka.plugins.event_producer.KafkaEventProducerPlugin attribute)
O
offset (airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload attribute)
on_dag_run_failed() (airflow.providers.apache.kafka.plugins.event_producer.DagRunListener method)
on_dag_run_running() (airflow.providers.apache.kafka.plugins.event_producer.DagRunListener method)
on_dag_run_success() (airflow.providers.apache.kafka.plugins.event_producer.DagRunListener method)
on_task_instance_running() (airflow.providers.apache.kafka.plugins.event_producer.TaskListener method)
open_stream() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer method)
P
partition (airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload attribute)
poll_interval (airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
poll_timeout (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger attribute)
producer_function (airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
producer_function() (in module tests.system.apache.kafka.example_dag_hello_kafka)
producer_function_args (airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
producer_function_kwargs (airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
ProduceToTopicOperator (class in airflow.providers.apache.kafka.operators.produce)
Q
queue_matches() (airflow.providers.apache.kafka.queues.kafka.KafkaMessageQueueProvider method)
QUEUE_REGEXP (in module airflow.providers.apache.kafka.queues.kafka)
R
read_to_end (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
return_apply_function_results (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
run() (airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger method)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger method)
S
sanitize_uri() (in module airflow.providers.apache.kafka.assets.kafka)
SCHEMA_VERSION (in module airflow.providers.apache.kafka.plugins.event_producer)
scheme (airflow.providers.apache.kafka.queues.kafka.KafkaMessageQueueProvider attribute)
serialize() (airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger method)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger method)
shared_stream_key() (airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger method)
synchronous (airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
T
t0 (in module tests.system.apache.kafka.example_dag_event_listener)
(in module tests.system.apache.kafka.example_dag_hello_kafka)
TaskListener (class in airflow.providers.apache.kafka.plugins.event_producer)
template_fields (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
test_connection() (airflow.providers.apache.kafka.hooks.base.KafkaBaseHook method)
test_run (in module tests.system.apache.kafka.example_dag_event_listener)
(in module tests.system.apache.kafka.example_dag_hello_kafka)
(in module tests.system.apache.kafka.example_dag_kafka_message_queue_trigger)
(in module tests.system.apache.kafka.example_dag_message_queue_trigger)
tests.system.apache.kafka
module
tests.system.apache.kafka.example_dag_event_listener
module
tests.system.apache.kafka.example_dag_hello_kafka
module
tests.system.apache.kafka.example_dag_kafka_message_queue_trigger
module
tests.system.apache.kafka.example_dag_message_queue_trigger
module
topic (airflow.providers.apache.kafka.operators.produce.ProduceToTopicOperator attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload attribute)
topics (airflow.providers.apache.kafka.hooks.consume.KafkaConsumerHook attribute)
(airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
(airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamProducer attribute)
(airflow.providers.apache.kafka.triggers.shared_stream.KafkaSharedStreamTrigger attribute)
trigger (in module tests.system.apache.kafka.example_dag_kafka_message_queue_trigger)
(in module tests.system.apache.kafka.example_dag_message_queue_trigger)
trigger_class() (airflow.providers.apache.kafka.queues.kafka.KafkaMessageQueueProvider method)
trigger_kwargs() (airflow.providers.apache.kafka.queues.kafka.KafkaMessageQueueProvider method)
U
ui_color (airflow.providers.apache.kafka.operators.consume.ConsumeFromTopicOperator attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
(airflow.providers.apache.kafka.sensors.kafka.AwaitMessageTriggerFunctionSensor attribute)
V
VALID_COMMIT_CADENCE (in module airflow.providers.apache.kafka.operators.consume)
(in module airflow.providers.apache.kafka.sensors.kafka)
value (airflow.providers.apache.kafka.triggers.shared_stream.KafkaBrokerPayload attribute)
X
xcom_push_key (airflow.providers.apache.kafka.sensors.kafka.AwaitMessageSensor attribute)
Previous
Next
Was this entry helpful?
Suggest a change on this page