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

44 statements  

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

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