为何第一个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, )
关键修复点说明
- 修正缩进问题:将
producer.send放入for循环内,确保每条消息处理完成后立即发送聚合数据。 - 优化Kafka配置:
- 给消费者添加
consumer_timeout_ms,避免无限阻塞导致Airflow任务超时。 - 使用
auto_offset_reset='earliest'确保消费到所有未处理的消息。 - 配置
value_deserializer和value_serializer,简化消息序列化/反序列化逻辑。
- 给消费者添加
- 添加异常处理:用
try/finally块确保消费者和生产者在任何情况下都能正常关闭,避免资源泄漏。 - 移除冗余代码:删除DAG末尾的
process_data调用,避免在解析阶段提前执行函数。
内容的提问来源于stack exchange,提问作者DEEPIKA MUTHUKUMAR
相关产品推荐
相关产品推荐

