Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/database/_migrations/versions/postgresql/2022_04_21_095519_14dc68cc5853_backfill_state_name.py: 92%

20 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

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 """