Apache Airflow绘图任务失败:Negsignal.SIGABRT与僵尸作业求助
问题
使用Apache Airflow读取CSV文件并生成可视化图表时,读取任务执行成功,但绘图任务失败,报错返回码Negsignal.SIGABRT且检测到僵尸作业。已确认内存充足,需修复问题使任务正常运行。
原代码如下:
from datetime import datetime from airflow import DAG from airflow.operators.python_operator import PythonOperator import pandas as pd import matplotlib.pyplot as plt import seaborn as sns import logging # Define function to read the CSV file def read_insurance_csv(**kwargs): try: # Path to the CSV file file_path = '/Users/xxx/airflow/dags/insurance_2.csv' # Read the CSV file insurance_df = pd.read_csv(file_path) # Save the DataFrame to a CSV file for verification insurance_df.to_csv('/Users/xxx/airflow/dags/insurance_read.csv', index=False) # Push the DataFrame to XCom kwargs['ti'].xcom_push(key='insurance_df', value=insurance_df.to_json()) # Log the DataFrame's head (first few rows) to verify the content logging.info("Successfully read insurance.csv and pushed to XCom.") print(insurance_df.head()) except Exception as e: logging.error(f"Error reading insurance.csv: {e}") raise # Define function to create plots def create_plots(**kwargs): try: # Pull the DataFrame from XCom ti = kwargs['ti'] insurance_json = ti.xcom_pull(task_ids='read_insurance_csv', key='insurance_df') df = pd.read_json(insurance_json) # Save the DataFrame to a CSV file for verification df.to_csv('/Users/xxx/airflow/dags/insurance_for_plotting.csv', index=False) # Define the plotting code numerical_features = ['age', 'bmi'] numerical_discrete = ['children'] categorical_features = ['sex', 'smoker', 'region'] # Create a figure and axes fig, axes = plt.subplots(2, 3, figsize=(15, 10)) # Iterate over numerical features for i, num_feature in enumerate(numerical_features): # Scatter plots for numerical vs numerical sns.scatterplot(x=num_feature, y='charges', data=df, ax=axes[0, i]) axes[0, i].set_title(f'Charges vs {num_feature}') for i, num_feature in enumerate(numerical_discrete): # Scatter plots for numerical vs numerical sns.boxplot(x=num_feature, y='charges', data=df, ax=axes[0, 2]) axes[0, 2].set_title(f'Charges vs {num_feature}') # Iterate over categorical features for i, cat_feature in enumerate(categorical_features): # Kernel density plots for numerical vs categorical sns.kdeplot(x='charges', hue=cat_feature, data=df, common_norm=False, ax=axes[1, i]) axes[1, i].set_title(f'Charges by {cat_feature}') plt.tight_layout() plt.savefig('/Users/xxx/airflow/dags/insurance_plots.png') plt.show() logging.info("Successfully created and saved plots.") except Exception as e: logging.error(f"Error creating plots: {e}") raise # Define default arguments for the DAG default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 4, 4), 'retries': 1, } # Define the DAG dag = DAG( 'read_and_plot_insurance_csv_dag', default_args=default_args, description='A DAG for reading the insurance CSV file and creating plots', schedule_interval=None, # Set to None for manual triggering ) # Define the task to read the CSV file read_insurance_csv_task = PythonOperator( task_id='read_insurance_csv', python_callable=read_insurance_csv, provide_context=True, dag=dag, ) # Define the task to create plots create_plots_task = PythonOperator( task_id='create_plots', python_callable=create_plots, provide_context=True, dag=dag, ) # Set task dependencies read_insurance_csv_task >> create_plots_task
原因分析
报错Negsignal.SIGABRT和僵尸作业的核心原因是:Airflow的Worker运行在无头(Headless)环境中,没有图形界面支持,但代码中调用了plt.show()——这个方法会尝试启动GUI窗口渲染图形,在无GUI环境下会触发进程异常终止,进而产生SIGABRT信号和僵尸作业。
解决方案
针对上述问题,只需做两处关键修改:
- 强制matplotlib使用非交互式后端(如
Agg),避免尝试启动GUI - 移除
plt.show()调用,因为Airflow环境不需要显示图形
修改后的代码
重点修改部分已标注:
from datetime import datetime from airflow import DAG from airflow.operators.python_operator import PythonOperator import pandas as pd import matplotlib.pyplot as plt import seaborn as sns import logging # -------------------------- # 修改1:设置matplotlib为非交互式后端 plt.switch_backend('Agg') # -------------------------- # Define function to read the CSV file def read_insurance_csv(**kwargs): try: file_path = '/Users/xxx/airflow/dags/insurance_2.csv' insurance_df = pd.read_csv(file_path) insurance_df.to_csv('/Users/xxx/airflow/dags/insurance_read.csv', index=False) kwargs['ti'].xcom_push(key='insurance_df', value=insurance_df.to_json()) logging.info("Successfully read insurance.csv and pushed to XCom.") print(insurance_df.head()) except Exception as e: logging.error(f"Error reading insurance.csv: {e}") raise # Define function to create plots def create_plots(**kwargs): try: ti = kwargs['ti'] insurance_json = ti.xcom_pull(task_ids='read_insurance_csv', key='insurance_df') df = pd.read_json(insurance_json) df.to_csv('/Users/xxx/airflow/dags/insurance_for_plotting.csv', index=False) numerical_features = ['age', 'bmi'] numerical_discrete = ['children'] categorical_features = ['sex', 'smoker', 'region'] fig, axes = plt.subplots(2, 3, figsize=(15, 10)) for i, num_feature in enumerate(numerical_features): sns.scatterplot(x=num_feature, y='charges', data=df, ax=axes[0, i]) axes[0, i].set_title(f'Charges vs {num_feature}') for i, num_feature in enumerate(numerical_discrete): sns.boxplot(x=num_feature, y='charges', data=df, ax=axes[0, 2]) axes[0, 2].set_title(f'Charges vs {num_feature}') for i, cat_feature in enumerate(categorical_features): sns.kdeplot(x='charges', hue=cat_feature, data=df, common_norm=False, ax=axes[1, i]) axes[1, i].set_title(f'Charges by {cat_feature}') plt.tight_layout() plt.savefig('/Users/xxx/airflow/dags/insurance_plots.png') # -------------------------- # 修改2:移除plt.show()调用 # plt.show() # -------------------------- logging.info("Successfully created and saved plots.") except Exception as e: logging.error(f"Error creating plots: {e}") raise default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 4, 4), 'retries': 1, } dag = DAG( 'read_and_plot_insurance_csv_dag', default_args=default_args, description='A DAG for reading the insurance CSV file and creating plots', schedule_interval=None, ) read_insurance_csv_task = PythonOperator( task_id='read_insurance_csv', python_callable=read_insurance_csv, provide_context=True, dag=dag, ) create_plots_task = PythonOperator( task_id='create_plots', python_callable=create_plots, provide_context=True, dag=dag, ) read_insurance_csv_task >> create_plots_task
额外优化建议
如果CSV文件数据量较大,XCom传递JSON序列化的DataFrame可能存在性能问题,可以改为:
- 读取任务直接将DataFrame保存为Parquet文件(比CSV更高效)
- 绘图任务直接读取Parquet文件,而非通过XCom传递数据
这样可以减少序列化/反序列化的开销,同时避免XCom的大小限制。
内容的提问来源于stack exchange,提问作者statsbeginner
相关产品推荐
相关产品推荐

