1"""Add state_timestamp
2
3Revision ID: 22b7cb02e593
4Revises: 2d5e000696f1
5Create Date: 2022-10-12 10:20:48.760447
6
7"""
8
9import sqlalchemy as sa
10from alembic import op
11
12import prefect
13
14# revision identifiers, used by Alembic.
15revision = "22b7cb02e593"
16down_revision = "2d5e000696f1"
17branch_labels = None
18depends_on = None
19
20
21def upgrade():
22 with op.batch_alter_table("flow_run", schema=None) as batch_op:
23 batch_op.add_column(
24 sa.Column(
25 "state_timestamp",
26 prefect.server.utilities.database.Timestamp(timezone=True),
27 nullable=True,
28 )
29 )
30 batch_op.create_index(
31 "ix_flow_run__state_timestamp", ["state_timestamp"], unique=False
32 )
33
34 with op.batch_alter_table("task_run", schema=None) as batch_op:
35 batch_op.add_column(
36 sa.Column(
37 "state_timestamp",
38 prefect.server.utilities.database.Timestamp(timezone=True),
39 nullable=True,
40 )
41 )
42 batch_op.create_index(
43 "ix_task_run__state_timestamp", ["state_timestamp"], unique=False
44 )
45
46 update_flow_run_state_timestamp_in_batches = """
47 WITH null_flow_run_state_timestamp_cte as (SELECT id from flow_run where state_timestamp is null and state_id is not null limit 500)
48 UPDATE flow_run
49 SET state_timestamp = flow_run_state.timestamp
50 FROM flow_run_state, null_flow_run_state_timestamp_cte
51 WHERE flow_run.state_id = flow_run_state.id
52 AND flow_run.id = null_flow_run_state_timestamp_cte.id;
53 """
54
55 update_task_run_state_timestamp_in_batches = """
56 WITH null_task_run_state_timestamp_cte as (SELECT id from task_run where state_timestamp is null and state_id is not null limit 500)
57 UPDATE task_run
58 SET state_timestamp = task_run_state.timestamp
59 FROM task_run_state, null_task_run_state_timestamp_cte
60 WHERE task_run.state_id = task_run_state.id
61 AND task_run.id = null_task_run_state_timestamp_cte.id;
62 """
63
64 with op.get_context().autocommit_block():
65 conn = op.get_bind()
66 while True:
67 # execute until we've backfilled all flow run state timestamps
68 # autocommit mode will commit each time `execute` is called
69 result = conn.execute(sa.text(update_flow_run_state_timestamp_in_batches))
70 if result.rowcount <= 0: 70 ↛ 66line 70 didn't jump to line 66 because the condition on line 70 was always true
71 break
72
73 while True:
74 # execute until we've backfilled all task run state timestamps
75 # autocommit mode will commit each time `execute` is called
76 result = conn.execute(sa.text(update_task_run_state_timestamp_in_batches))
77 if result.rowcount <= 0: 77 ↛ 73line 77 didn't jump to line 73 because the condition on line 77 was always true
78 break
79
80
81def downgrade():
82 with op.batch_alter_table("task_run", schema=None) as batch_op:
83 batch_op.drop_index("ix_task_run__state_timestamp")
84 batch_op.drop_column("state_timestamp")
85
86 with op.batch_alter_table("flow_run", schema=None) as batch_op:
87 batch_op.drop_index("ix_flow_run__state_timestamp")
88 batch_op.drop_column("state_timestamp")