Airflow子任务出错但任务未失败的原因排查
Great that you've already pinpointed the root schema issue (dw2 doesn't exist), let's dive into why Airflow isn't flagging your computer_unload_and_load task as failed despite the error showing up in logs.
The Core Issue in Your Code
The problem lies in how you're handling the subprocess.run() call and exception handling in your load_table function:
subprocess.run()doesn't fail by default: When yourpsqlcommand fails (because the schema is missing), it returns a non-zero exit code—butsubprocess.run()won't throw an exception unless you explicitly tell it to. Without that trigger, your Python function just runs to completion, returnsNone, and Airflow interprets that as a successful task run.- Your try/except block isn't catching the failure: The try/except only catches Python-level exceptions, but since
subprocess.run()isn't raising an error for the failed psql command, there's nothing to catch or re-raise. Airflow never gets the signal that something went wrong.
Confirmation From Your Logs
Look at the final line in your task logs:
[2018-05-30 11:22:46,969] {base_task_runner.py:98} INFO - Subtask: [2018-05-30 11:22:46,968] {python_operator.py:90} INFO - Done. Returned value was: None
This confirms the function finished executing without any uncaught exceptions. Airflow's PythonOperator marks a task as successful if the callable completes without raising an error—regardless of whether the underlying command (like psql) failed.
Quick Fix to Ensure Failure is Detected
Since you already know how to fix the schema issue, here's how to adjust your code so Airflow properly marks the task as failed when the psql command fails:
Either use the check=True parameter with subprocess.run() to force an exception on non-zero exit codes, or manually check the return code:
def load_table(host, db, schema, table, unload_task_id=False, file_path=False, **kwargs): """ load a csv file into a table if no file_path is given, uses XCOM to get the file_name returned by the unload task""" try: if not file_path: file_path = kwargs['ti'].xcom_pull(task_ids=unload_task_id) load_cmd = "\copy {}.{} FROM {} WITH (FORMAT CSV, NULL '^', HEADER)".format(schema, table, file_path) command = [ "psql", "-U", "root", "-h", host, "-d", db, "-c", load_cmd ] # Enable check=True to raise an exception if psql fails subprocess.run(command, check=True) except Exception as e: # Optional: Add logging here for better debugging print(f"Load operation failed: {str(e)}") raise # Re-raise the exception so Airflow marks the task as failed
内容的提问来源于stack exchange,提问作者smullan

