# ใช้ SQL สำหรับแก้ปัญหาตอนรัน Airflow Migration หลังจากอัพเกรดเวอร์ชั่น

**URL:** <https://discuss.dataengineercafe.io/t/sql-airflow-migration/11>\
**Category:** Airflow\
**Tags:** airflow, sql, migration\
**Created:** [November 11, 2021, 10:16am UTC](https://discuss.dataengineercafe.io/t/sql-airflow-migration/11 "2021-11-11T10:16:02Z")\
**Posts on this page:** 1\
**Page:** 1

<div class="post-metadata">

**Author:** ![zkan](https://yyz1.discourse-cdn.com/flex035/user_avatar/discuss.dataengineercafe.io/zkan/32/2_2.png) [@zkan](https://discuss.dataengineercafe.io/u/zkan)\
**Post date:** [November 11, 2021, 10:16am UTC](https://discuss.dataengineercafe.io/t/sql-airflow-migration/11/1 "2021-11-11T10:16:02Z")

</div>

ตอนที่เปลี่ยนมาเป็น Airflow เวอร์ชั่น 2.2.0 แล้ว ตัว migration รันแล้วเจอ error แบบนี้

```auto
[2021-10-11 19:21:13,069] {db.py:817} ERROR - The task_instance table has 53 rows without a corresponding dag_run row. You must manually correct this problem (possibly by deleting the problem rows).
[2021-10-11 19:21:13,071] {db.py:817} ERROR - The task_fail table has 24 rows without a corresponding dag_run row. You must manually correct this problem (possibly by deleting the problem rows).

```

เหมือนมีข้อมูลที่ไม่ได้ลิ้งค์กับใครเลยค้างอยู่ เราต้องเคลียร์พวกนี้ออกก่อน วิธีแก้ปัญหาตอนนี้คือใช้ SQL ใน issue ด้านล่างนี้เลย

```sql
BEGIN;

-- Remove dag runs without a valid run_id
DELETE FROM dag_run WHERE run_id is NULL;

-- Remove task fails without a run_id
WITH task_fails_to_remove AS (
  SELECT 
    task_fail.dag_id,
    task_fail.task_id,
    task_fail.execution_date
  FROM
    task_fail
  LEFT JOIN 
    dag_run ON 
    dag_run.dag_id = task_fail.dag_id 
    AND dag_run.execution_date = task_fail.execution_date
  WHERE
    dag_run.run_id IS NULL
)
DELETE FROM
    task_fail
USING 
    task_fails_to_remove
WHERE (
    task_fail.dag_id = task_fails_to_remove.dag_id
    AND task_fail.task_id = task_fails_to_remove.task_id
    AND task_fail.execution_date = task_fails_to_remove.execution_date
);

-- Remove task instances without a run_id
WITH task_instances_to_remove AS (
  SELECT
    task_instance.dag_id,
    task_instance.task_id,
    task_instance.execution_date
  FROM
    task_instance
  LEFT JOIN 
    dag_run 
    ON dag_run.dag_id = task_instance.dag_id
    AND dag_run.execution_date = task_instance.execution_date
  WHERE 
    dag_run.run_id is NULL
)
DELETE FROM 
    task_instance
USING
    task_instances_to_remove
WHERE (
    task_instance.dag_id = task_instances_to_remove.dag_id
    AND task_instance.task_id = task_instances_to_remove.task_id
    AND task_instance.execution_date = task_instances_to_remove.execution_date
);

COMMIT;

```

> <https://github.com/apache/airflow/issues/18894#issuecomment-941325572>
>
> \### Apache Airflow version
> 
> 2.2.0
> 
> \### Operating System
> 
> Linux
> 
> \### Vers…ions of Apache Airflow Providers
> 
> default.
> 
> \### Deployment
> 
> Docker-Compose
> 
> \### Deployment details
> 
> Using airflow-2.2.0python3.7
> 
> \### What happened
> 
> Upgrading image from apache/airflow:2.1.4-python3.7
> to apache/airflow:2.2.0-python3.7
> Cause this inside scheduler, which is not starting:
> 
> \`\`\`
> Python version: 3.7.12
> Airflow version: 2.2.0
> Node: 6dd55b0a5dd7
> \-------------------------------------------------------------------------------
> Traceback (most recent call last):
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1277, in \_execute\_context
> cursor, statement, parameters, context
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/default.py", line 608, in do\_execute
> cursor.execute(statement, parameters)
> psycopg2.errors.UndefinedColumn: column dag.max\_active\_tasks does not exist
> LINE 1: ..., dag.schedule\_interval AS dag\_schedule\_interval, dag.max\_ac...
> ^
> 
> The above exception was the direct cause of the following exception:
> 
> Traceback (most recent call last):
> File "/home/airflow/.local/lib/python3.7/site-packages/flask/app.py", line 2447, in wsgi\_app
> response = self.full\_dispatch\_request()
> File "/home/airflow/.local/lib/python3.7/site-packages/flask/app.py", line 1952, in full\_dispatch\_request
> rv = self.handle\_user\_exception(e)
> File "/home/airflow/.local/lib/python3.7/site-packages/flask/app.py", line 1821, in handle\_user\_exception
> reraise(exc\_type, exc\_value, tb)
> File "/home/airflow/.local/lib/python3.7/site-packages/flask/\_compat.py", line 39, in reraise
> raise value
> File "/home/airflow/.local/lib/python3.7/site-packages/flask/app.py", line 1950, in full\_dispatch\_request
> rv = self.dispatch\_request()
> File "/home/airflow/.local/lib/python3.7/site-packages/flask/app.py", line 1936, in dispatch\_request
> return self.view\_functions\[rule.endpoint\](\*\*req.view\_args)
> File "/home/airflow/.local/lib/python3.7/site-packages/airflow/www/auth.py", line 51, in decorated
> return func(\*args, \*\*kwargs)
> File "/home/airflow/.local/lib/python3.7/site-packages/airflow/www/views.py", line 588, in index
> filter\_dag\_ids = current\_app.appbuilder.sm.get\_accessible\_dag\_ids(g.user)
> File "/home/airflow/.local/lib/python3.7/site-packages/airflow/www/security.py", line 377, in get\_accessible\_dag\_ids
> return {dag.dag\_id for dag in accessible\_dags}
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/orm/query.py", line 3535, in \_\_iter\_\_
> return self.\_execute\_and\_instances(context)
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/orm/query.py", line 3560, in \_execute\_and\_instances
> result = conn.execute(querycontext.statement, self.\_params)
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1011, in execute
> return meth(self, multiparams, params)
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/sql/elements.py", line 298, in \_execute\_on\_connection
> return connection.\_execute\_clauseelement(self, multiparams, params)
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1130, in \_execute\_clauseelement
> distilled\_params,
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1317, in \_execute\_context
> e, statement, parameters, cursor, context
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1511, in \_handle\_dbapi\_exception
> sqlalchemy\_exception, with\_traceback=exc\_info\[2\], from\_=e
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/util/compat.py", line 182, in raise\_
> raise exception
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1277, in \_execute\_context
> cursor, statement, parameters, context
> File "/home/airflow/.local/lib/python3.7/site-packages/sqlalchemy/engine/default.py", line 608, in do\_execute
> cursor.execute(statement, parameters)
> sqlalchemy.exc.ProgrammingError: (psycopg2.errors.UndefinedColumn) column dag.max\_active\_tasks does not exist
> LINE 1: ..., dag.schedule\_interval AS dag\_schedule\_interval, dag.max\_ac...
> ^
> 
> \[SQL: SELECT dag.dag\_id AS dag\_dag\_id, dag.root\_dag\_id AS dag\_root\_dag\_id, dag.is\_paused AS dag\_is\_paused, dag.is\_subdag AS dag\_is\_subdag, dag.is\_active AS dag\_is\_active, dag.last\_parsed\_time AS dag\_last\_parsed\_time, dag.last\_pickled AS dag\_last\_pickled, dag.last\_expired AS dag\_last\_expired, dag.scheduler\_lock AS dag\_scheduler\_lock, dag.pickle\_id AS dag\_pickle\_id, dag.fileloc AS dag\_fileloc, dag.owners AS dag\_owners, dag.description AS dag\_description, dag.default\_view AS dag\_default\_view, dag.schedule\_interval AS dag\_schedule\_interval, dag.max\_active\_tasks AS dag\_max\_active\_tasks, dag.max\_active\_runs AS dag\_max\_active\_runs, dag.has\_task\_concurrency\_limits AS dag\_has\_task\_concurrency\_limits, dag.next\_dagrun AS dag\_next\_dagrun, dag.next\_dagrun\_data\_interval\_start AS dag\_next\_dagrun\_data\_interval\_start, dag.next\_dagrun\_data\_interval\_end AS dag\_next\_dagrun\_data\_interval\_end, dag.next\_dagrun\_create\_after AS dag\_next\_dagrun\_create\_after 
> FROM dag\]
> (Background on this error at: http://sqlalche.me/e/13/f405)
> \`\`\`
> 
> \### What you expected to happen
> 
> Automatic database migration and properly working scheduler.
> 
> \### How to reproduce
> 
> Ugrade from 2.1.4 to 2.2.0 with some dags history.
> 
> \### Anything else
> 
> \_No response\_
> 
> \### Are you willing to submit PR?
> 
> \- \[\] Yes I am willing to submit a PR!
> 
> \### Code of Conduct
> 
> \- \[X\] I agree to follow this project's \[Code of Conduct\](https://github.com/apache/airflow/blob/main/CODE\_OF\_CONDUCT.md)

ล่าสุดใช้ Query นี้ก็น่าจะเพียงพอ

```sql
WITH task_fails_to_remove AS (
  SELECT 
    task_fail.dag_id,
    task_fail.task_id,
    task_fail.execution_date
  FROM
    task_fail
  LEFT JOIN 
    dag_run ON 
    dag_run.dag_id = task_fail.dag_id 
    AND dag_run.execution_date = task_fail.execution_date
  WHERE
    dag_run.run_id IS NULL
)
DELETE FROM
    task_fail
USING 
    task_fails_to_remove
WHERE (
    task_fail.dag_id = task_fails_to_remove.dag_id
    AND task_fail.task_id = task_fails_to_remove.task_id
    AND task_fail.execution_date = task_fails_to_remove.execution_date
);

```

ปัญหานี้เดี๋ยวทาง Airflow Community จะแก้ให้เวลาที่เราลบ DAG Run มันจะลบแบบ cascade (ไปลบตัวอื่นๆ เช่น task instance ที่เกี่ยวข้องด้วย)
