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

22 statements  

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

1"""Migrate artifact data to artifact_collection table 

2 

3Revision ID: 15f5083c16bd 

4Revises: 310dda75f561 

5Create Date: 2023-04-06 12:27:16.676260 

6 

7""" 

8 

9import sqlalchemy as sa 

10from alembic import op 

11 

12# revision identifiers, used by Alembic. 

13revision = "15f5083c16bd" 

14down_revision = "310dda75f561" 

15branch_labels = None 

16depends_on = None 

17 

18 

19def upgrade(): 

20 """ 

21 A data-only migration that populates flow_run_id, task_run_id, type, description, and metadata_ columns 

22 for artifact_collection table. 

23 """ 

24 batch_size = 500 

25 offset = 0 

26 

27 update_artifact_collection_table = """ 

28 WITH artifact_collection_cte AS ( 

29 SELECT * FROM artifact_collection WHERE id = :id FOR UPDATE 

30 ) 

31 UPDATE artifact_collection 

32 SET data = artifact.data, 

33 description = artifact.description, 

34 flow_run_id = artifact.flow_run_id, 

35 task_run_id = artifact.task_run_id, 

36 type = artifact.type, 

37 metadata_ = artifact.metadata_ 

38 FROM artifact, artifact_collection_cte 

39 WHERE artifact_collection.latest_id = artifact.id 

40 AND artifact.id = artifact_collection_cte.latest_id; 

41 """ 

42 

43 with op.get_context().autocommit_block(): 

44 conn = op.get_bind() 

45 while True: 

46 select_artifact_collection_cte = f""" 

47 SELECT * from artifact_collection ORDER BY id LIMIT {batch_size} OFFSET {offset} FOR UPDATE; 

48 """ 

49 

50 # Get the next batch of rows to update 

51 selected_artifact_collections = conn.execute( 

52 sa.text(select_artifact_collection_cte) 

53 ).fetchall() 

54 if not selected_artifact_collections: 54 ↛ 57line 54 didn't jump to line 57 because the condition on line 54 was always true

55 break 

56 

57 for row in selected_artifact_collections: 

58 id_to_update = row[0] 

59 conn.execute( 

60 sa.text(update_artifact_collection_table), {"id": id_to_update} 

61 ) 

62 offset += batch_size 

63 

64 

65def downgrade(): 

66 """ 

67 Data-only migration, no action needed. 

68 """