20260731_08a_create_video_discovery_tables.py 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277
  1. """create the complete video discovery persistence tables
  2. Revision ID: 20260731_08a
  3. Revises: 20260730_07
  4. Create Date: 2026-07-31
  5. The tables predate Alembic and were previously installed from files under
  6. ``sql/``. ``checkfirst``-style inspection keeps this migration safe for
  7. databases where those legacy DDL files were already applied, while making a
  8. fresh ``alembic upgrade head`` self-contained.
  9. """
  10. from __future__ import annotations
  11. from collections.abc import Sequence
  12. import sqlalchemy as sa
  13. from alembic import op
  14. revision: str = "20260731_08a"
  15. down_revision: str | None = "20260730_07"
  16. branch_labels: str | Sequence[str] | None = None
  17. depends_on: str | Sequence[str] | None = None
  18. def _table_names() -> set[str]:
  19. return set(sa.inspect(op.get_bind()).get_table_names())
  20. def _create_run_table() -> None:
  21. op.create_table(
  22. "video_discovery_run",
  23. sa.Column("id", sa.BigInteger(), autoincrement=True, nullable=False),
  24. sa.Column("run_id", sa.String(length=64), nullable=False, comment="Agent 运行标识"),
  25. sa.Column("biz_dt", sa.String(length=32), nullable=True, comment="业务日 YYYYMMDD"),
  26. sa.Column("demand_grade_id", sa.BigInteger(), nullable=True),
  27. sa.Column("demand_word", sa.String(length=256), nullable=False),
  28. sa.Column("seed_video_id", sa.String(length=64), nullable=True),
  29. sa.Column("seed_video_title", sa.String(length=512), nullable=True),
  30. sa.Column("relevant_points_json", sa.Text(), nullable=False),
  31. sa.Column("intent_summary", sa.Text(), nullable=True),
  32. sa.Column(
  33. "status",
  34. sa.String(length=24),
  35. nullable=False,
  36. server_default="running",
  37. ),
  38. sa.Column("search_count", sa.Integer(), nullable=False, server_default="0"),
  39. sa.Column("primary_count", sa.Integer(), nullable=False, server_default="0"),
  40. sa.Column("stop_reason", sa.Text(), nullable=True),
  41. sa.Column(
  42. "create_time",
  43. sa.DateTime(),
  44. nullable=False,
  45. server_default=sa.func.now(),
  46. ),
  47. sa.Column(
  48. "update_time",
  49. sa.DateTime(),
  50. nullable=False,
  51. server_default=sa.func.now(),
  52. ),
  53. sa.PrimaryKeyConstraint("id"),
  54. sa.UniqueConstraint("run_id", name="uk_video_discovery_run_id"),
  55. sa.UniqueConstraint(
  56. "biz_dt",
  57. "demand_grade_id",
  58. name="uk_video_discovery_run_biz_grade",
  59. ),
  60. mysql_charset="utf8mb4",
  61. )
  62. op.create_index(
  63. "idx_video_discovery_run_demand",
  64. "video_discovery_run",
  65. ["demand_word"],
  66. )
  67. op.create_index(
  68. "idx_video_discovery_run_grade",
  69. "video_discovery_run",
  70. ["demand_grade_id"],
  71. )
  72. op.create_index(
  73. "idx_video_discovery_run_biz_dt",
  74. "video_discovery_run",
  75. ["biz_dt"],
  76. )
  77. op.create_index(
  78. "idx_video_discovery_run_status",
  79. "video_discovery_run",
  80. ["status"],
  81. )
  82. def _create_search_table() -> None:
  83. op.create_table(
  84. "video_discovery_search",
  85. sa.Column("id", sa.BigInteger(), autoincrement=True, nullable=False),
  86. sa.Column("run_id", sa.String(length=64), nullable=False),
  87. sa.Column("search_key", sa.String(length=64), nullable=False),
  88. sa.Column("keyword", sa.String(length=256), nullable=False),
  89. sa.Column("query_reason", sa.Text(), nullable=False),
  90. sa.Column("source_type", sa.String(length=32), nullable=False),
  91. sa.Column("source_value", sa.Text(), nullable=True),
  92. sa.Column("parent_search_id", sa.BigInteger(), nullable=True),
  93. sa.Column(
  94. "provider",
  95. sa.String(length=32),
  96. nullable=False,
  97. server_default="internal_keyword",
  98. ),
  99. sa.Column("provider_state_json", sa.Text(), nullable=True),
  100. sa.Column(
  101. "content_type",
  102. sa.String(length=16),
  103. nullable=False,
  104. server_default="视频",
  105. ),
  106. sa.Column(
  107. "sort_type",
  108. sa.String(length=32),
  109. nullable=False,
  110. server_default="综合排序",
  111. ),
  112. sa.Column(
  113. "publish_time",
  114. sa.String(length=32),
  115. nullable=False,
  116. server_default="不限",
  117. ),
  118. sa.Column(
  119. "cursor",
  120. sa.String(length=128),
  121. nullable=False,
  122. server_default="0",
  123. ),
  124. sa.Column("page_no", sa.Integer(), nullable=False, server_default="1"),
  125. sa.Column("results_count", sa.Integer(), nullable=False, server_default="0"),
  126. sa.Column(
  127. "new_candidate_count",
  128. sa.Integer(),
  129. nullable=False,
  130. server_default="0",
  131. ),
  132. sa.Column("has_more", sa.Integer(), nullable=False, server_default="0"),
  133. sa.Column("next_cursor", sa.String(length=128), nullable=True),
  134. sa.Column("result_ids_json", sa.Text(), nullable=True),
  135. sa.Column(
  136. "status",
  137. sa.String(length=16),
  138. nullable=False,
  139. server_default="success",
  140. ),
  141. sa.Column("error_message", sa.Text(), nullable=True),
  142. sa.Column(
  143. "create_time",
  144. sa.DateTime(),
  145. nullable=False,
  146. server_default=sa.func.now(),
  147. ),
  148. sa.Column(
  149. "update_time",
  150. sa.DateTime(),
  151. nullable=False,
  152. server_default=sa.func.now(),
  153. ),
  154. sa.PrimaryKeyConstraint("id"),
  155. mysql_charset="utf8mb4",
  156. )
  157. op.create_index(
  158. "idx_video_discovery_search_run",
  159. "video_discovery_search",
  160. ["run_id", "id"],
  161. )
  162. op.create_index(
  163. "idx_video_discovery_search_parent",
  164. "video_discovery_search",
  165. ["run_id", "parent_search_id"],
  166. )
  167. op.create_index(
  168. "idx_video_discovery_search_keyword",
  169. "video_discovery_search",
  170. ["keyword"],
  171. )
  172. def _create_candidate_table() -> None:
  173. op.create_table(
  174. "video_discovery_candidate",
  175. sa.Column("id", sa.BigInteger(), autoincrement=True, nullable=False),
  176. sa.Column("run_id", sa.String(length=64), nullable=False),
  177. sa.Column(
  178. "search_id",
  179. sa.BigInteger(),
  180. sa.ForeignKey(
  181. "video_discovery_search.id",
  182. name="fk_video_discovery_candidate_search",
  183. ondelete="RESTRICT",
  184. ),
  185. nullable=True,
  186. ),
  187. sa.Column("aweme_id", sa.String(length=64), nullable=False),
  188. sa.Column("title", sa.String(length=512), nullable=True),
  189. sa.Column("content_link", sa.String(length=1024), nullable=True),
  190. sa.Column("author_name", sa.String(length=256), nullable=True),
  191. sa.Column("author_sec_uid", sa.String(length=256), nullable=True),
  192. sa.Column("source_keywords_json", sa.Text(), nullable=True),
  193. sa.Column("source_search_ids_json", sa.Text(), nullable=True),
  194. sa.Column("tags_json", sa.Text(), nullable=True),
  195. sa.Column("play_count", sa.BigInteger(), nullable=True),
  196. sa.Column("like_count", sa.BigInteger(), nullable=True),
  197. sa.Column("comment_count", sa.BigInteger(), nullable=True),
  198. sa.Column("collect_count", sa.BigInteger(), nullable=True),
  199. sa.Column("share_count", sa.BigInteger(), nullable=True),
  200. sa.Column("content_age_evidence_json", sa.Text(), nullable=True),
  201. sa.Column("account_age_evidence_json", sa.Text(), nullable=True),
  202. sa.Column("age_normalization_json", sa.Text(), nullable=True),
  203. sa.Column("relevance_score", sa.Numeric(8, 6), nullable=True),
  204. sa.Column("elder_score", sa.Numeric(8, 6), nullable=True),
  205. sa.Column("share_score", sa.Numeric(8, 6), nullable=True),
  206. sa.Column("value_score", sa.Numeric(8, 2), nullable=True),
  207. sa.Column("decision_reason", sa.Text(), nullable=True),
  208. sa.Column(
  209. "decision_bucket",
  210. sa.String(length=24),
  211. nullable=False,
  212. server_default="pending_evaluation",
  213. ),
  214. sa.Column("aigc_crawler_plan_id", sa.String(length=64), nullable=True),
  215. sa.Column("aigc_produce_plan_id", sa.String(length=64), nullable=True),
  216. sa.Column("aigc_publish_plan_id", sa.String(length=64), nullable=True),
  217. sa.Column("aigc_plan_label", sa.String(length=64), nullable=True),
  218. sa.Column(
  219. "create_time",
  220. sa.DateTime(),
  221. nullable=False,
  222. server_default=sa.func.now(),
  223. ),
  224. sa.Column(
  225. "update_time",
  226. sa.DateTime(),
  227. nullable=False,
  228. server_default=sa.func.now(),
  229. ),
  230. sa.PrimaryKeyConstraint("id"),
  231. mysql_charset="utf8mb4",
  232. )
  233. op.create_index(
  234. "idx_video_discovery_candidate_search",
  235. "video_discovery_candidate",
  236. ["search_id", "id"],
  237. )
  238. op.create_index(
  239. "idx_video_discovery_candidate_bucket",
  240. "video_discovery_candidate",
  241. ["run_id", "decision_bucket"],
  242. )
  243. op.create_index(
  244. "idx_video_discovery_candidate_author",
  245. "video_discovery_candidate",
  246. ["author_sec_uid"],
  247. )
  248. def upgrade() -> None:
  249. tables = _table_names()
  250. if "video_discovery_run" not in tables:
  251. _create_run_table()
  252. if "video_discovery_search" not in tables:
  253. _create_search_table()
  254. if "video_discovery_candidate" not in tables:
  255. _create_candidate_table()
  256. def downgrade() -> None:
  257. # These tables may have existed before Alembic took ownership of the schema.
  258. # A destructive automatic downgrade could delete legacy production data, so
  259. # this adoption migration is intentionally forward-only.
  260. pass