feat: 球队实体一致性 — 归一化咽喉 + team_aliases 别名机制

events/standings 创建 Team 前均经 team_names.normalize(已有,确认),
TeamRepository.get_or_create 收敛为归一化唯一咽喉 + info 日志。

新增 team_aliases 表(NFKD 归一别名 → teams.id FK CASCADE),
定位三步链:normalize(name) → teams.name → team_aliases → insert。
不自动合并历史重复队;提供 POST /api/v1/admin/teams/aliases 显式添加。

迁移 0020_team_aliases + Admin 别名管理端点(admin_teams.py)。
全量测试 270 通过。
This commit is contained in:
shangfangjian
2026-09-22 00:27:53 +08:00
parent 4b0d6ee58a
commit e15b554ba3
10 changed files with 290 additions and 15 deletions
+32
View File
@@ -0,0 +1,32 @@
"""球队别名表 team_aliases
Revision ID: 0020_team_aliases
Revises: 0019_ingest_jobs
Create Date: 2026-09-22
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = '0020_team_aliases'
down_revision: Union[str, None] = '0019_ingest_jobs'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.create_table(
'team_aliases',
sa.Column('alias_normalized', sa.String(120), primary_key=True),
sa.Column('team_id', sa.Integer, sa.ForeignKey('teams.id', ondelete='CASCADE'), nullable=False),
sa.Column('original_alias', sa.String(120), nullable=False),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
)
op.create_index('ix_team_aliases_team_id', 'team_aliases', ['team_id'])
def downgrade() -> None:
op.drop_index('ix_team_aliases_team_id', table_name='team_aliases')
op.drop_table('team_aliases')
+27 -1
View File
@@ -60,7 +60,7 @@
### 队名归一化 ### 队名归一化
`src/data/team_names.py` 维护 `NORMALIZE_MAP`(如 `Man City` → `Manchester City`),未命中映射的队名原样返回。 `src/data.team_names.py` 维护 `NORMALIZE_MAP`(如 `Man City` → `Manchester City`),未命中映射的队名原样返回。
归一前先做 Unicode NFKD 去重音。 归一前先做 Unicode NFKD 去重音。
**唯一键是归一后英文名**:`teams.name` 带 `UNIQUE` 约束,所有入库路径均经 `TeamRepository.get_or_create` 收敛归一化 **唯一键是归一后英文名**:`teams.name` 带 `UNIQUE` 约束,所有入库路径均经 `TeamRepository.get_or_create` 收敛归一化
@@ -70,6 +70,32 @@
>(如 `"Man City"` → `"Manchester City"`,但 `"man city"` 原样保留)。上游 bzzoiro 返回的队名首字母大写, >(如 `"Man City"` → `"Manchester City"`,但 `"man city"` 原样保留)。上游 bzzoiro 返回的队名首字母大写,
>实际命中无问题;若新增数据源返回全小写/全大写队名,需先 `title()` 再归一,否则会绕过映射产生重复 Team。 >实际命中无问题;若新增数据源返回全小写/全大写队名,需先 `title()` 再归一,否则会绕过映射产生重复 Team。
#### 别名机制(`team_aliases`)
归一仍可能遗漏历史重复队(如 `"Bayern Munich"` 与 `"Bayern München"` 经 NFKD 后相同则命中,
但 `"Man United"` vs `"Manchester United"` 若漏映射)。`team_aliases` 表提供**显式别名→teams.id** 映射:
| 列 | 说明 |
|----|------|
| `alias_normalized` | PK,`normalize(别名)` 后的稳定幂等键 |
| `team_id` | FK → `teams.id`(ON DELETE CASCADE) |
| `original_alias` | 原始写法(保留供参考) |
**定位三步链**(`get_or_create`):`normalize(name)` → 查 `teams.name` → 查 `team_aliases`(以 `normalize(name)` 为 PK)→ 都没有才 insert 新 Team。别名命中即复用已有 Team,避免产生重复。
**添加别名**(不自动合并历史重复队):
- **Admin 接口**(推荐):`POST /api/v1/admin/teams/aliases {"alias": "Man United", "team_id": 42}`(require_admin,幂等)
- **直接 SQL**:
```sql
INSERT INTO team_aliases(alias_normalized, team_id, original_alias)
VALUES ('man united', 42, 'Man United')
ON CONFLICT (alias_normalized) DO UPDATE SET team_id = EXCLUDED.team_id, original_alias = EXCLUDED.original_alias;
```
> ⚠️ **别名不自动合并**:发现历史重复队 A/B 后,需人工确认归一目标(如保留 B),再为 A 的归一名添加别名指向 B。
> 合并前请确认 A 的 `matches`/`standings` 引用是否需要迁移(可先 `SELECT COUNT(*) FROM matches WHERE home_team_id = A.id OR away_team_id = A.id` 评估)。
**改名 / 合并流程**(人工): **改名 / 合并流程**(人工):
当发现两个 `teams` 行实际是同一球队(如 `Manchester City` 与 `Man City` 因历史数据大小写差异各占一行): 当发现两个 `teams` 行实际是同一球队(如 `Manchester City` 与 `Man City` 因历史数据大小写差异各占一行):
+4
View File
@@ -15,11 +15,15 @@ from fastapi import APIRouter
from src.api.routes.admin_config import router as admin_config_router from src.api.routes.admin_config import router as admin_config_router
from src.api.routes.admin_datasources import router as admin_datasources_router from src.api.routes.admin_datasources import router as admin_datasources_router
from src.api.routes.admin_ingest_jobs import router as admin_ingest_jobs_router
from src.api.routes.admin_llm import router as admin_llm_router from src.api.routes.admin_llm import router as admin_llm_router
from src.api.routes.admin_quality import router as admin_quality_router from src.api.routes.admin_quality import router as admin_quality_router
from src.api.routes.admin_teams import router as admin_teams_router
router = APIRouter() router = APIRouter()
router.include_router(admin_datasources_router) router.include_router(admin_datasources_router)
router.include_router(admin_config_router) router.include_router(admin_config_router)
router.include_router(admin_ingest_jobs_router)
router.include_router(admin_llm_router) router.include_router(admin_llm_router)
router.include_router(admin_quality_router) router.include_router(admin_quality_router)
router.include_router(admin_teams_router)
+60
View File
@@ -0,0 +1,60 @@
"""后台管理:球队别名管理(只读列表 + 添加别名)。
归一名(teams.name)是球队唯一键;别名(team_aliases)是同一球队的不同写法
(大小写/译名/缩写)到归一后 teams.id 的映射。入库时 normalize(name) 依次查
teams.name 与 team_aliases,命中即复用,避免重复 Team。
不自动合并历史重复队;需显式添加别名(或先 SQL/再经由此接口)。
"""
from __future__ import annotations
import logging
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy import desc, select
from src.api.deps import require_admin
from src.api.schemas import TeamAliasIn, TeamAliasOut
from src.db.base import AsyncSession, get_db_read
from src.db.models import Team, TeamAlias
from src.db.repositories import TeamRepository
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/v1/admin", tags=["admin"], dependencies=[Depends(require_admin)])
@router.get("/teams/aliases", response_model=list[TeamAliasOut])
async def list_team_aliases(db: AsyncSession = Depends(get_db_read)):
"""列出所有球队别名(最新在前)。"""
rows = (await db.execute(select(TeamAlias).order_by(desc(TeamAlias.created_at)).limit(200))).scalars().all()
return [
TeamAliasOut(
alias_normalized=r.alias_normalized,
team_id=r.team_id,
original_alias=r.original_alias,
)
for r in rows
]
@router.post("/teams/aliases", response_model=TeamAliasOut, status_code=201)
async def add_team_alias(req: TeamAliasIn, db: AsyncSession = Depends(get_db_read)):
"""为已有 Team 添加别名(幂等:重复添加会更新指向)。
不自动合并历史重复队。若需合并 A→B:先为 A 的归一名添加别名指向 B,
再人工确认 A 是否仍有独立引用。
"""
# 校验目标 Team 存在
team = await db.get(Team, req.team_id)
if team is None:
raise HTTPException(404, f"目标 Team 不存在: id={req.team_id}")
repo = TeamRepository(db)
row = await repo.add_alias(req.alias, req.team_id)
logger.info("添加 Team 别名: %s -> team_id=%s", req.alias, req.team_id)
return TeamAliasOut(
alias_normalized=row.alias_normalized,
team_id=row.team_id,
original_alias=row.original_alias,
)
+35
View File
@@ -117,6 +117,19 @@ class IngestBzzoiroRequest(BaseModel):
season: str | None = Field(None, description="standings 赛季,如 '2026-2027';空 = 当前赛季") season: str | None = Field(None, description="standings 赛季,如 '2026-2027';空 = 当前赛季")
class TeamAliasIn(BaseModel):
"""POST /api/v1/admin/teams/aliases 请求体:为已有 Team 添加别名。"""
alias: str = Field(..., min_length=1, max_length=120, description="球队别名(原始写法)")
team_id: int = Field(..., gt=0, description="归一后的目标 teams.id")
class TeamAliasOut(BaseModel):
alias_normalized: str
team_id: int
original_alias: str
class IngestResponse(BaseModel): class IngestResponse(BaseModel):
leagues: dict leagues: dict
total_inserted: int total_inserted: int
@@ -124,6 +137,28 @@ class IngestResponse(BaseModel):
errors: list[str] = [] errors: list[str] = []
class IngestBzzoiroResponse(BaseModel):
"""POST /api/v1/ingest/bzzoiro 响应:兼容原 message 字段,新增 job_id 供轮询。"""
ok: bool = True
job_id: str = Field(..., description="采集任务 ID(GET /api/v1/admin/ingest/jobs/{job_id} 轮询)")
message: str = ""
class IngestJobOut(BaseModel):
"""采集任务状态详情。"""
id: str
task: str
params: dict
status: str # pending | running | success | failed
result: dict | None = None
error: str | None = None
created_at: datetime | None = None
started_at: datetime | None = None
finished_at: datetime | None = None
class ScheduleIn(BaseModel): class ScheduleIn(BaseModel):
id: str = Field(..., description="任务唯一标识,如 'daily-events'") id: str = Field(..., description="任务唯一标识,如 'daily-events'")
task: str = Field(..., description="events / standings / stats / all") task: str = Field(..., description="events / standings / stats / all")
+45
View File
@@ -56,6 +56,22 @@ class Team(Base):
away_matches: Mapped[list["Match"]] = relationship(foreign_keys="Match.away_team_id", back_populates="away_team") away_matches: Mapped[list["Match"]] = relationship(foreign_keys="Match.away_team_id", back_populates="away_team")
class TeamAlias(Base):
"""球队别名:同一球队的不同写法(大小写/译名/缩写)映射到归一后的 teams.id。
入库流程(get_or_create):normalize(name) → 查 teams.name → 查 team_aliases
→ 都没有再 insert 新 Team。别名不自动合并历史重复队,需显式添加。
alias_normalized 为 normalize(别名)后的稳定幂等键,用作 PK 避免重复插入。
"""
__tablename__ = "team_aliases"
# normalize(别名)后的值,稳定幂等,用作主键
alias_normalized: Mapped[str] = mapped_column(String(120), primary_key=True)
team_id: Mapped[int] = mapped_column(ForeignKey("teams.id", ondelete="CASCADE"), nullable=False)
original_alias: Mapped[str] = mapped_column(String(120), nullable=False) # 原始写法(保留供参考)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_utcnow)
class Match(Base): class Match(Base):
__tablename__ = "matches" __tablename__ = "matches"
@@ -324,6 +340,35 @@ class RawEvent(Base):
) )
class IngestJob(Base):
"""采集任务状态:跟踪每次触后台采集任务的执行进度与结果。
POST /api/v1/ingest/bzzoiro 触发时写入(pending→running→success/failed),
前端 Collection 页据此轮询到终态,替代此前"30 秒后盲标完成"的模拟。
分批 get_uow / BzzoiroSource / IngestFailure / Bronze/Lineage 均不受影响
(本表仅作状态追踪,不介入采集事务)。
"""
__tablename__ = "ingest_jobs"
id: Mapped[str] = mapped_column(String(36), primary_key=True) # uuid4
task: Mapped[str] = mapped_column(String(20), nullable=False)
params: Mapped[dict] = mapped_column(JSONB, nullable=False, default=dict)
status: Mapped[str] = mapped_column(String(20), nullable=False, server_default="pending")
result: Mapped[dict | None] = mapped_column(JSONB)
error: Mapped[str | None] = mapped_column(Text)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
__table_args__ = (
Index("ix_ingest_job_status_created", "status", "created_at"),
CheckConstraint(
"status IN ('pending', 'running', 'success', 'failed')",
name="ck_ingest_job_status",
),
)
class IngestFailure(Base): class IngestFailure(Base):
"""采集失败死信:记录失败原因、重试次数与下次重试时间。 """采集失败死信:记录失败原因、重试次数与下次重试时间。
+42 -7
View File
@@ -114,24 +114,59 @@ class TeamRepository:
async def get_or_create(self, name: str, *, name_zh: str | None = None) -> Team: async def get_or_create(self, name: str, *, name_zh: str | None = None) -> Team:
"""按名获取球队,不存在则创建(name_zh 供 bzzoiro 管线写中文名)。 """按名获取球队,不存在则创建(name_zh 供 bzzoiro 管线写中文名)。
归一化咽喉:所有入库 Team.name 必须经过 team_names.normalize, 归一化咽喉 + 别名查找,三步定位:
此处统一收敛,避免各调用点散落归一化逻辑导致重复 Team。 1) normalize(name) → 查 teams.name
2) 查 team_aliases(以 normalize(name) 为幂等键)→ 复用已映射的 teams.id
3) 都没有 → insert 新 Team(归一名)
创建新 Team 时 info 打出原始名与归一后的规范名,便于排查重名。 创建新 Team 时 info 打出原始名与归一后的规范名,便于排查重名。
不自动合并历史重复队;需显式添加别名。
""" """
from src.data.team_names import normalize as normalize_name from src.data.team_names import normalize as normalize_name
from src.db.models import TeamAlias
normalized = normalize_name(name) or name.strip() normalized = normalize_name(name) or name.strip()
# 1) 归一名直查 teams
team = await self.get_by_name(normalized) team = await self.get_by_name(normalized)
if team is None: if team is not None:
logger.info( return team
"创建新 Team: %s -> %s",
name, normalized, # 2) 别名查找:normalize(别名) 作为幂等键,命中即复用已有 Team
) alias = await self._session.get(TeamAlias, normalized)
if alias is not None:
team = await self._session.get(Team, alias.team_id)
if team is not None:
logger.info("Team 别名命中: %s -> %s(已有 id=%s)", name, normalized, team.id)
return team
# 3) 新建 Team(归一名)
logger.info("创建新 Team: %s -> %s", name, normalized)
team = Team(name=normalized, name_zh=name_zh) team = Team(name=normalized, name_zh=name_zh)
self._session.add(team) self._session.add(team)
await self._session.flush() await self._session.flush()
return team return team
async def add_alias(self, alias: str, team_id: int) -> TeamAlias:
"""为已有 Team 添加别名。
幂等:以 normalize(alias) 为 PK,重复添加同一别名会 upsert。
不自动合并历史重复队,仅建立别名映射。
"""
from src.data.team_names import normalize as normalize_name
from src.db.models import TeamAlias
normalized = normalize_name(alias) or alias.strip()
existing = await self._session.get(TeamAlias, normalized)
if existing is not None:
existing.team_id = team_id # 允许重新指向
existing.original_alias = alias
await self._session.flush()
return existing
row = TeamAlias(alias_normalized=normalized, team_id=team_id, original_alias=alias)
self._session.add(row)
await self._session.flush()
return row
async def get_all_by_names(self, names: list[str]) -> dict[str, Team]: async def get_all_by_names(self, names: list[str]) -> dict[str, Team]:
"""批量获取球队,返回 name → Team 映射。""" """批量获取球队,返回 name → Team 映射。"""
if not names: if not names:
+13 -2
View File
@@ -20,7 +20,7 @@ from datetime import date, datetime, timezone
import pytest import pytest
import src.data.bzzoiro as bz import src.data.bzzoiro as bz
from src.db.models import DataLineage, League, Match, RawEvent, Team from src.db.models import DataLineage, League, Match, RawEvent, Team, TeamAlias
def _event(eid=1001, status="finished", home="Arsenal", away="Chelsea", hs=2, as_=1): def _event(eid=1001, status="finished", home="Arsenal", away="Chelsea", hs=2, as_=1):
@@ -71,11 +71,20 @@ class _FakeDB:
League: list(leagues), League: list(leagues),
RawEvent: list(raw_events), RawEvent: list(raw_events),
} }
self._next_id = 0 self._teams_by_id: dict[int, Team] = {t.id: t for t in teams if getattr(t, "id", None)}
self._aliases: dict[str, TeamAlias] = {}
self._next_id = max((t.id for t in teams if getattr(t, "id", None)), default=0)
def add(self, obj): def add(self, obj):
self.added.append(obj) self.added.append(obj)
async def get(self, cls, key):
if cls is Team:
return self._teams_by_id.get(key)
if cls is TeamAlias:
return self._aliases.get(key)
return None
async def execute(self, stmt): async def execute(self, stmt):
entities = set() entities = set()
for d in (stmt.column_descriptions or []): for d in (stmt.column_descriptions or []):
@@ -90,6 +99,8 @@ class _FakeDB:
if getattr(obj, "id", None) is None: if getattr(obj, "id", None) is None:
self._next_id += 1 self._next_id += 1
obj.id = self._next_id obj.id = self._next_id
if isinstance(obj, Team) and obj.id is not None:
self._teams_by_id[obj.id] = obj
@pytest.fixture(autouse=True) @pytest.fixture(autouse=True)
+14
View File
@@ -17,6 +17,7 @@ import re
import pytest import pytest
from src.db.models import Team, TeamAlias
from src.data.key_ring import _mask from src.data.key_ring import _mask
from src.llm import backtest as bt_mod from src.llm import backtest as bt_mod
from src.llm.agents import orchestrator as orch_mod from src.llm.agents import orchestrator as orch_mod
@@ -105,6 +106,15 @@ class _FakeDb:
self.added: list = [] self.added: list = []
self.flush_count = 0 self.flush_count = 0
self._next_id = 1000 self._next_id = 1000
self._teams_by_id: dict[int, Team] = {}
self._aliases: dict[str, TeamAlias] = {}
async def get(self, cls, key):
if cls is Team:
return self._teams_by_id.get(key)
if cls is TeamAlias:
return self._aliases.get(key)
return None
async def execute(self, _stmt): async def execute(self, _stmt):
if self._results: if self._results:
@@ -120,6 +130,10 @@ class _FakeDb:
if getattr(obj, "id", None) is None: if getattr(obj, "id", None) is None:
self._next_id += 1 self._next_id += 1
obj.id = self._next_id obj.id = self._next_id
if isinstance(obj, Team) and getattr(obj, "id", None) is not None:
self._teams_by_id[obj.id] = obj
if isinstance(obj, TeamAlias):
self._aliases[obj.alias_normalized] = obj
async def test_r2_standings_actually_upserts(monkeypatch): async def test_r2_standings_actually_upserts(monkeypatch):
+15 -2
View File
@@ -18,7 +18,7 @@ from __future__ import annotations
import pytest import pytest
import src.data.bzzoiro as bz import src.data.bzzoiro as bz
from src.db.models import DataLineage, IngestFailure, League, RawEvent, Standing, Team from src.db.models import DataLineage, IngestFailure, League, RawEvent, Standing, Team, TeamAlias
def _payload(): def _payload():
@@ -78,11 +78,21 @@ class _FakeDB:
Standing: list(standings), Standing: list(standings),
RawEvent: list(raw_events), RawEvent: list(raw_events),
} }
self._next_id = 0 # session.get 查找表(Team/TeamAlias)
self._teams_by_id: dict[int, Team] = {t.id: t for t in teams if getattr(t, "id", None)}
self._aliases: dict[str, TeamAlias] = {}
self._next_id = max((t.id for t in teams if getattr(t, "id", None)), default=0)
def add(self, obj): def add(self, obj):
self.added.append(obj) self.added.append(obj)
async def get(self, cls, key):
if cls is Team:
return self._teams_by_id.get(key)
if cls is TeamAlias:
return self._aliases.get(key)
return None
async def execute(self, stmt): async def execute(self, stmt):
entities = set() entities = set()
for d in (stmt.column_descriptions or []): for d in (stmt.column_descriptions or []):
@@ -110,6 +120,9 @@ class _FakeDB:
if getattr(obj, "id", None) is None: if getattr(obj, "id", None) is None:
self._next_id += 1 self._next_id += 1
obj.id = self._next_id obj.id = self._next_id
# 同步 session.get 可查到新建 Team
if isinstance(obj, Team) and obj.id is not None:
self._teams_by_id[obj.id] = obj
def _preset_league(): def _preset_league():