You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow中无法在日志查看Kafka Topic消息的问题排查

问题解决:Airflow Kafka消费任务日志未打印消息

核心问题分析

你的代码中func函数存在逻辑错误,导致无法正确处理Kafka消息并输出到日志:

  1. func的参数是单个Kafka消息,但你错误地循环遍历get_messages(这是Airflow任务实例,而非消息集合)
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 08:56:05