Airflow中无法在日志查看Kafka Topic消息的问题排查
问题解决:Airflow Kafka消费任务日志未打印消息
核心问题分析
你的代码中func函数存在逻辑错误,导致无法正确处理Kafka消息并输出到日志:
func的参数是单个Kafka消息,但你错误地循环遍历get_messages(这是Airflow任务实例,而非消息集合)ConsumeFromTopicOperator会将每条消息单独传递给apply_function,不需要额外循环遍历
修正后的代码
from airflow import DAG from datetime import datetime, timedelta from airflow_provider_kafka.operators.consume_from_topic import ConsumeFromTopicOperator import logging def func(message, prefix=None): # 直接处理传入的单条Kafka消息 message_value = str(message.value) # 使用Airflow日志模块输出(比print更规范,确保日志被Airflow捕获) logging.info(f"Kafka消息内容: {message_value}") # 若坚持用print,也会被写入日志,但logging更适配Airflow生态 # print(f"Kafka消息内容: {message_value}") with DAG( dag_id="test_kafka", start_date=datetime(2021, 1, 1), schedule_interval='@weekly', catchup=False # 建议添加,避免启动后立即执行历史调度任务 ) as dag: get_messages = ConsumeFromTopicOperator( task_id="get_messages", topics=["topictest"], apply_function='test_kafka.func', consumer_config={ 'group.id': 'test-consumer-group', 'bootstrap.servers': 'server:9092', "auto.offset.reset": "earliest", # 可选:若消息量少,添加参数确保消费者能获取到消息 # 'fetch.min.bytes': 1, # 'fetch.wait.max.ms': 500 } ) get_messages
额外排查项
- 确认Airflow Worker节点能够访问Kafka集群地址
server:9092 - 检查
topictest主题是否存在且包含有效消息 - 若之前用同一消费者组消费过,可更换
group.id或重置该组的offset,确保能重新拉取消息 - 验证Airflow Worker的日志级别设置为INFO及以上(默认是INFO,若改为WARNING则不会显示info级别的日志)
内容的提问来源于stack exchange,提问作者VLBigo
相关产品推荐
相关产品推荐

