Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/database/_migrations/versions/postgresql/2023_01_26_045501_2882cd2df464_create_migration_index.py: 82%

34 statements  

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

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