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

为何第一个DAG触发第二个DAG后,后者发送SIGTERM信号并失败?

问题排查与修复:Airflow DAG执行Kafka消费时触发SIGTERM终止

你的第二个DAG出现SIGTERM终止的问题,主要由代码中的缩进错误、消费者配置缺失以及冗余代码导致,以下是具体问题分析和修复方案:

核心问题分析

  • 缩进错误导致变量未定义:producer.send语句缩进错误,既不在for循环内,也不属于process_data函数,执行时会找不到aggregated_json变量,直接触发异常终止任务。
  • Kafka消费者无超时配置:默认的KafkaConsumer会一直阻塞等待新消息,Airflow任务可能因超时而发送SIGTERM信号终止进程。
  • 资源清理逻辑不健壮:未通过try/finally确保消费者和生产者在异常情况下也能正常关闭,可能引发资源泄漏。
  • 冗余代码触发提前执行:DAG末尾的process_data会在DAG解析阶段直接调用函数,而非任务运行时执行,可能导致不必要的错误。

修复后的完整代码

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from kafka.producer import KafkaProducer
from kafka.consumer.group import KafkaConsumer
import json

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2023, 6, 16),
    'retries': 0,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    'data_processor_dag',
    default_args=default_args,
    schedule_interval='@once'
)

def process_data():
    # 配置Kafka消费者,添加超时与偏移量设置
    consumer = KafkaConsumer(
        'checkdata',
        bootstrap_servers='localhost:9092',
        auto_offset_reset='earliest',  # 从最早的消息开始消费
        consumer_timeout_ms=10000,     # 10秒无消息则停止消费
        value_deserializer=lambda m: json.loads(m.decode('utf-8'))  # 自动反序列化消息
    )

    # 配置Kafka生产者,自动序列化JSON
    producer = KafkaProducer(
        bootstrap_servers='localhost:9092',
        value_serializer=lambda m: json.dumps(m).encode('utf-8')
    )

    try:
        # 遍历消费消息
        for message in consumer:
            data = message.value
            message_data = data['event']['messageData']
            timestamp1 = data['timestamp']
            nanos = data['nanos']
            vin = data['vin']

            # 计算聚合指标
            Safetyratings = [msg['SafetyRating'] for msg in message_data]
            PerformanceRatings = [msg['PerformanceRating'] for msg in message_data]
            EfficiencyRatings = [msg['EfficiencyRating'] for msg in message_data]
            distance = [msg['distance_travelled'] for msg in message_data]
            
            actual_distance = round(sum(distance), 0)
            avg_safetyrating = round(sum(Safetyratings) / len(Safetyratings), 0)
            avg_performancerating = round(sum(PerformanceRatings) / len(PerformanceRatings), 0)
            avg_efficiencyrating = round(sum(EfficiencyRatings) / len(EfficiencyRatings), 0)
            Score = round(((avg_efficiencyrating + avg_performancerating + avg_safetyrating)/3), 0)

            # 构造聚合数据并发送
            aggregated_json = {
                "SafetyRating": avg_safetyrating,
                "PerformanceRating": avg_performancerating,
                "EfficiencyRating": avg_efficiencyrating,
                "Distance_travelled": actual_distance,
                "Score": Score,
                "vin": vin,
                "timestamp": timestamp1,
                "nanos": nanos
            }
            producer.send('aggregated_data', aggregated_json)
            
        # 确保所有消息发送完成
        producer.flush()
        print("Data processing complete")
    finally:
        # 无论是否异常,都关闭资源
        consumer.close()
        producer.close()

process_data_task = PythonOperator(
    task_id='process_data',
    python_callable=process_data,
    dag=dag,
)

关键修复点说明

  1. 修正缩进问题:将producer.send放入for循环内,确保每条消息处理完成后立即发送聚合数据。
  2. 优化Kafka配置:
    • 给消费者添加consumer_timeout_ms,避免无限阻塞导致Airflow任务超时。
    • 使用auto_offset_reset='earliest'确保消费到所有未处理的消息。
    • 配置value_deserializer和value_serializer,简化消息序列化/反序列化逻辑。
  3. 添加异常处理:用try/finally块确保消费者和生产者在任何情况下都能正常关闭,避免资源泄漏。
  4. 移除冗余代码:删除DAG末尾的process_data调用,避免在解析阶段提前执行函数。

内容的提问来源于stack exchange,提问作者DEEPIKA MUTHUKUMAR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:43:08