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