1"""Backfill state_name
2
3Revision ID: 14dc68cc5853
4Revises: 605ebb4e9155
5Create Date: 2022-04-21 09:55:19.820177
6
7"""
8
9import sqlalchemy as sa
10from alembic import op
11
12# revision identifiers, used by Alembic.
13revision = "14dc68cc5853"
14down_revision = "605ebb4e9155"
15branch_labels = None
16depends_on = None
17
18
19def upgrade():
20 """
21 Backfills state_name column for task_run and flow_run tables.
22
23 This is a data only migration that can be run as many
24 times as desired.
25 """
26
27 update_flow_run_state_name_in_batches = """
28 WITH null_flow_run_state_name_cte as (SELECT id from flow_run where state_name is null and state_id is not null limit 500)
29 UPDATE flow_run
30 SET state_name = flow_run_state.name
31 FROM flow_run_state, null_flow_run_state_name_cte
32 WHERE flow_run.state_id = flow_run_state.id
33 AND flow_run.id = null_flow_run_state_name_cte.id;
34 """
35
36 update_task_run_state_name_in_batches = """
37 WITH null_task_run_state_name_cte as (SELECT id from task_run where state_name is null and state_id is not null limit 500)
38 UPDATE task_run
39 SET state_name = task_run_state.name
40 FROM task_run_state, null_task_run_state_name_cte
41 WHERE task_run.state_id = task_run_state.id
42 AND task_run.id = null_task_run_state_name_cte.id;
43 """
44
45 with op.get_context().autocommit_block():
46 conn = op.get_bind()
47 while True:
48 # execute until we've backfilled all flow run state names
49 # autocommit mode will commit each time `execute` is called
50 result = conn.execute(sa.text(update_flow_run_state_name_in_batches))
51 if result.rowcount <= 0: 51 ↛ 47line 51 didn't jump to line 47 because the condition on line 51 was always true
52 break
53
54 while True:
55 # execute until we've backfilled all task run state names
56 # autocommit mode will commit each time `execute` is called
57 result = conn.execute(sa.text(update_task_run_state_name_in_batches))
58 if result.rowcount <= 0: 58 ↛ 54line 58 didn't jump to line 54 because the condition on line 58 was always true
59 break
60
61
62def downgrade():
63 """
64 Data only migration. No action on downgrade.
65 """