1"""Adds a helper index for the artifact data migration
2
3Revision ID: 2882cd2df464
4Revises: 2882cd2df463
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 = "2882cd2df464"
14down_revision = "2882cd2df463"
15branch_labels = None
16depends_on = None
17
18
19def upgrade():
20 with op.batch_alter_table("flow_run_state", schema=None) as batch_op:
21 batch_op.add_column(sa.Column("has_data", sa.Boolean))
22 batch_op.create_index(
23 batch_op.f("ix_flow_run_state__has_data"),
24 ["has_data"],
25 unique=False,
26 )
27
28 with op.batch_alter_table("task_run_state", schema=None) as batch_op:
29 batch_op.add_column(sa.Column("has_data", sa.Boolean))
30 batch_op.create_index(
31 batch_op.f("ix_task_run_state__has_data"),
32 ["has_data"],
33 unique=False,
34 )
35
36 def populate_flow_has_data_in_batches(batch_size):
37 return f"""
38 UPDATE flow_run_state
39 SET has_data = (data IS NOT NULL AND data != 'null')
40 WHERE flow_run_state.id in (SELECT id FROM flow_run_state WHERE (has_data IS NULL) LIMIT {batch_size});
41 """
42
43 def populate_task_has_data_in_batches(batch_size):
44 return f"""
45 UPDATE task_run_state
46 SET has_data = (data IS NOT NULL AND data != 'null')
47 WHERE task_run_state.id in (SELECT id FROM task_run_state WHERE (has_data IS NULL) LIMIT {batch_size});
48 """
49
50 migration_statements = [
51 populate_flow_has_data_in_batches,
52 populate_task_has_data_in_batches,
53 ]
54
55 with op.get_context().autocommit_block():
56 conn = op.get_bind()
57 for query in migration_statements:
58 batch_size = 500
59
60 while True:
61 # execute until we've updated task_run_state_id and artifact_data
62 # autocommit mode will commit each time `execute` is called
63 sql_stmt = sa.text(query(batch_size))
64 result = conn.execute(sql_stmt)
65
66 if result.rowcount < batch_size: 66 ↛ 60line 66 didn't jump to line 60 because the condition on line 66 was always true
67 break
68
69
70def downgrade():
71 with op.batch_alter_table("task_run_state", schema=None) as batch_op:
72 batch_op.drop_index(batch_op.f("ix_task_run_state__has_data"))
73 batch_op.drop_column("has_data")
74
75 with op.batch_alter_table("flow_run_state", schema=None) as batch_op:
76 batch_op.drop_index(batch_op.f("ix_flow_run_state__has_data"))
77 batch_op.drop_column("has_data")