20260728_01_drop_unused_scheduler_artifacts.py 2.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475
  1. """drop unused scheduler leftovers
  2. Revision ID: 20260728_01
  3. Revises: 20260727_03
  4. Create Date: 2026-07-28
  5. """
  6. from __future__ import annotations
  7. from collections.abc import Sequence
  8. import sqlalchemy as sa
  9. from alembic import op
  10. from sqlalchemy.dialects import mysql
  11. revision: str = "20260728_01"
  12. down_revision: str | None = "20260727_03"
  13. branch_labels: str | Sequence[str] | None = None
  14. depends_on: str | Sequence[str] | None = None
  15. def upgrade() -> None:
  16. op.execute("DROP TABLE IF EXISTS scheduler_job_execution")
  17. op.drop_column("pipeline_run", "triggered_by")
  18. op.drop_column("pipeline_run", "parent_run_id")
  19. op.drop_column("pipeline_run", "scheduled_for")
  20. op.drop_column("pipeline_step_run", "metrics_json")
  21. def downgrade() -> None:
  22. op.add_column(
  23. "pipeline_step_run",
  24. sa.Column("metrics_json", sa.JSON(), nullable=True),
  25. )
  26. op.add_column(
  27. "pipeline_run",
  28. sa.Column("scheduled_for", sa.DateTime(), nullable=True),
  29. )
  30. op.add_column(
  31. "pipeline_run",
  32. sa.Column("parent_run_id", sa.String(length=36), nullable=True),
  33. )
  34. op.add_column(
  35. "pipeline_run",
  36. sa.Column("triggered_by", sa.String(length=128), nullable=True),
  37. )
  38. op.create_table(
  39. "scheduler_job_execution",
  40. sa.Column("id", sa.BigInteger(), autoincrement=True, nullable=False),
  41. sa.Column("run_id", sa.String(length=64), nullable=False),
  42. sa.Column("job_name", sa.String(length=128), nullable=False),
  43. sa.Column("job_id", sa.String(length=128), nullable=False),
  44. sa.Column("status", sa.String(length=32), nullable=False),
  45. sa.Column("event_time", sa.DateTime(), nullable=False),
  46. sa.Column("biz_dt", sa.String(length=8), nullable=True),
  47. sa.Column("started_at", sa.DateTime(), nullable=True),
  48. sa.Column("finished_at", sa.DateTime(), nullable=True),
  49. sa.Column("duration_seconds", sa.Float(), nullable=True),
  50. sa.Column("error_message", sa.Text(), nullable=True),
  51. sa.Column("detail", mysql.LONGTEXT(), nullable=True),
  52. sa.Column(
  53. "create_time",
  54. sa.DateTime(),
  55. server_default=sa.text("CURRENT_TIMESTAMP"),
  56. nullable=False,
  57. ),
  58. sa.Column(
  59. "update_time",
  60. sa.DateTime(),
  61. server_default=sa.text("CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP"),
  62. nullable=False,
  63. ),
  64. sa.PrimaryKeyConstraint("id"),
  65. mysql_charset="utf8mb4",
  66. mysql_collate="utf8mb4_unicode_ci",
  67. )