1"""Migrates state data to the artifact table
2
3Revision ID: 2882cd2df465
4Revises: 2882cd2df464
5Create Date: 2023-01-26 04:55:01.358638
6
7"""
8
9import sqlalchemy as sa
10from alembic import op
11
12# revision identifiers, used by Alembic.
13revision = "2882cd2df465"
14down_revision = "2882cd2df464"
15branch_labels = None
16depends_on = None
17
18
19def upgrade():
20 ### START DATA MIGRATION
21
22 # insert nontrivial task run state results into the artifact table
23 def update_task_run_artifact_data_in_batches(batch_size, offset):
24 return f"""
25 INSERT INTO artifact (task_run_state_id, task_run_id, data)
26 SELECT id, task_run_id, data
27 FROM task_run_state
28 WHERE has_data IS TRUE
29 ORDER BY id
30 LIMIT {batch_size} OFFSET {offset};
31 """
32
33 # backpopulate the result artifact id on the task run state table
34 def update_task_run_state_from_artifact_id_in_batches(batch_size, offset):
35 return f"""
36 UPDATE task_run_state
37 SET result_artifact_id = (SELECT id FROM artifact WHERE task_run_state.id = task_run_state_id)
38 WHERE task_run_state.id in (SELECT id FROM task_run_state WHERE (has_data IS TRUE) AND (result_artifact_id IS NULL) LIMIT {batch_size});
39 """
40
41 # insert nontrivial flow run state results into the artifact table
42 def update_flow_run_artifact_data_in_batches(batch_size, offset):
43 return f"""
44 INSERT INTO artifact (flow_run_state_id, flow_run_id, data)
45 SELECT id, flow_run_id, data
46 FROM flow_run_state
47 WHERE has_data IS TRUE
48 ORDER BY id
49 LIMIT {batch_size} OFFSET {offset};
50 """
51
52 # backpopulate the result artifact id on the flow run state table
53 def update_flow_run_state_from_artifact_id_in_batches(batch_size, offset):
54 return f"""
55 UPDATE flow_run_state
56 SET result_artifact_id = (SELECT id FROM artifact WHERE flow_run_state.id = flow_run_state_id)
57 WHERE flow_run_state.id in (SELECT id FROM flow_run_state WHERE (has_data IS TRUE) AND (result_artifact_id IS NULL) LIMIT {batch_size});
58 """
59
60 data_migration_queries = [
61 update_task_run_artifact_data_in_batches,
62 update_task_run_state_from_artifact_id_in_batches,
63 update_flow_run_artifact_data_in_batches,
64 update_flow_run_state_from_artifact_id_in_batches,
65 ]
66
67 with op.get_context().autocommit_block():
68 conn = op.get_bind()
69 for query in data_migration_queries:
70 batch_size = 500
71 offset = 0
72
73 while True:
74 # execute until we've updated task_run_state_id and artifact_data
75 # autocommit mode will commit each time `execute` is called
76 sql_stmt = sa.text(query(batch_size, offset))
77 result = conn.execute(sql_stmt)
78
79 if result.rowcount <= 0: 79 ↛ 82line 79 didn't jump to line 82 because the condition on line 79 was always true
80 break
81
82 offset += batch_size
83
84 ### END DATA MIGRATION
85
86
87def downgrade():
88 def nullify_artifact_ref_from_flow_run_state_in_batches(batch_size):
89 return f"""
90 UPDATE flow_run_state
91 SET result_artifact_id = NULL
92 WHERE flow_run_state.id in (SELECT id FROM flow_run_state WHERE result_artifact_id IS NOT NULL LIMIT {batch_size});
93 """
94
95 def nullify_artifact_ref_from_task_run_state_in_batches(batch_size):
96 return f"""
97 UPDATE task_run_state
98 SET result_artifact_id = NULL
99 WHERE task_run_state.id in (SELECT id FROM task_run_state WHERE result_artifact_id IS NOT NULL LIMIT {batch_size});
100 """
101
102 def delete_artifacts_in_batches(batch_size):
103 return f"""
104 DELETE FROM artifact
105 WHERE artifact.id IN (SELECT id FROM artifact LIMIT {batch_size});
106 """
107
108 data_migration_queries = [
109 delete_artifacts_in_batches,
110 nullify_artifact_ref_from_flow_run_state_in_batches,
111 nullify_artifact_ref_from_task_run_state_in_batches,
112 ]
113
114 with op.get_context().autocommit_block():
115 conn = op.get_bind()
116 for query in data_migration_queries:
117 batch_size = 500
118
119 while True:
120 # execute until we've updated task_run_state_id and artifact_data
121 # autocommit mode will commit each time `execute` is called
122 sql_stmt = sa.text(query(batch_size))
123 result = conn.execute(sql_stmt)
124
125 if result.rowcount <= 0:
126 break