feat(sector-radar): 提供版本化排名查询接口

This commit is contained in:
yuxuanhui
2026-08-29 19:44:46 +08:00
parent d9bae722d8
commit efc4c3d530
10 changed files with 1115 additions and 12 deletions
@@ -5,6 +5,7 @@ from fastapi import APIRouter
from zhixing_server.interfaces.http.system import operational_router, system_router
from zhixing_server.modules.market_data.presentation.home import home_router
from zhixing_server.modules.market_data.presentation.integrity import integrity_router
from zhixing_server.modules.sector_radar.presentation.http import sector_radar_router
from zhixing_server.modules.selection.presentation.http import selection_router
api_v1_router = APIRouter(prefix="/api/v1")
@@ -16,5 +17,10 @@ api_v1_router.include_router(
tags=["market-data"],
)
api_v1_router.include_router(selection_router, prefix="/selection", tags=["selection"])
api_v1_router.include_router(
sector_radar_router,
prefix="/sector-radar",
tags=["sector-radar"],
)
__all__ = ["api_v1_router", "operational_router"]
@@ -0,0 +1,223 @@
"""Stable read model for persisted sector radar publications."""
from __future__ import annotations
from dataclasses import dataclass
from datetime import date
from enum import StrEnum
from typing import Literal
from ..domain.metrics import (
AmountNetStrategy,
RatioTurnoverStrategy,
SwingEqualThreeToTenStrategy,
)
from ..domain.models import (
MetricKind,
MetricUnit,
RadarPublication,
RankedMetric,
RankSide,
SectorType,
)
from ..domain.persistence import SectorRadarRepository
from ..domain.ranking import select_percentile_side, select_rank_change_side
ReadStatus = Literal["success", "no_data"]
class RadarView(StrEnum):
"""Supported ranking projections at the HTTP boundary."""
AMOUNT = "amount"
RATIO = "ratio"
SWING = "swing"
RANK_CHANGE = "rank_change"
@dataclass(frozen=True, slots=True)
class RadarMetricDefinition:
"""Public definition of one explicitly independent metric implementation."""
metric_kind: MetricKind
metric_version: str
label: str
unit: MetricUnit
implementation_kind: Literal["independent"] = "independent"
disclaimer: str = "知行独立实现,非 OneChartLab 原站公式"
@dataclass(frozen=True, slots=True)
class RadarQuery:
"""Validated application query for one ranking page."""
trade_date: date | None = None
sector_type: SectorType = SectorType.CONCEPT
view: RadarView = RadarView.AMOUNT
rank_change_metric: MetricKind = MetricKind.AMOUNT
rank_change_days: int = 1
side: RankSide = RankSide.ALL
search: str | None = None
page: int = 1
page_size: int = 20
def __post_init__(self) -> None:
"""Reject invalid pagination and rank-history offsets outside HTTP usage."""
if not 1 <= self.rank_change_days <= 5:
raise ValueError("rank_change_days must be between 1 and 5")
if self.page < 1:
raise ValueError("page must be positive")
if not 1 <= self.page_size <= 100:
raise ValueError("page_size must be between 1 and 100")
if self.search is not None and len(self.search) > 100:
raise ValueError("search must not exceed 100 characters")
@dataclass(frozen=True, slots=True)
class RadarDateIndex:
"""Available successful dates plus the newest attempt and strict last-good."""
available_dates: tuple[date, ...]
current_attempt: RadarPublication | None
last_good: RadarPublication | None
@property
def status(self) -> ReadStatus:
"""Return no_data until at least one successful publication exists."""
return "success" if self.last_good is not None else "no_data"
@dataclass(frozen=True, slots=True)
class RankingPage:
"""One filtered page without losing publication or metric provenance."""
status: ReadStatus
query: RadarQuery
publication: RadarPublication | None
definition: RadarMetricDefinition
rows: tuple[RankedMetric, ...]
total: int
_METRIC_DEFINITIONS = {
MetricKind.AMOUNT: RadarMetricDefinition(
metric_kind=MetricKind.AMOUNT,
metric_version=AmountNetStrategy.metric_version,
label="主力净流入(知行独立实现)",
unit=MetricUnit.CNY_100M,
),
MetricKind.RATIO: RadarMetricDefinition(
metric_kind=MetricKind.RATIO,
metric_version=RatioTurnoverStrategy.metric_version,
label="主力净流入/成交额(知行独立实现)",
unit=MetricUnit.RATIO,
),
MetricKind.SWING: RadarMetricDefinition(
metric_kind=MetricKind.SWING,
metric_version=SwingEqualThreeToTenStrategy.metric_version,
label="3—10 日等权资金率(知行独立实现)",
unit=MetricUnit.RATIO,
),
}
class ReadSectorRadar:
"""Hide last-good selection, ranking filters, search, and pagination."""
def __init__(self, repository: SectorRadarRepository) -> None:
self.repository = repository
def list_dates(self) -> RadarDateIndex:
"""Return successful dates without promoting partial or failed attempts."""
return RadarDateIndex(
available_dates=tuple(self.repository.list_successful_dates()),
current_attempt=self.repository.get_latest_publication(),
last_good=self.repository.get_last_good_publication(),
)
def query(self, query: RadarQuery) -> RankingPage:
"""Return one deterministic page for ordinary or rank-change views."""
metric_kind = (
query.rank_change_metric
if query.view is RadarView.RANK_CHANGE
else MetricKind(query.view.value)
)
definition = _METRIC_DEFINITIONS[metric_kind]
publication = (
self.repository.get_successful_publication(query.trade_date)
if query.trade_date is not None
else self.repository.get_last_good_publication()
)
if publication is None:
return RankingPage("no_data", query, None, definition, (), 0)
metric_rows = tuple(
row
for row in self.repository.load_rankings(publication.publication_id)
if row.observation.sector_type is query.sector_type
and row.observation.metric_kind is metric_kind
and row.observation.metric_version == definition.metric_version
)
if query.view is RadarView.RANK_CHANGE:
if query.side is RankSide.ALL:
selected = tuple(
sorted(
metric_rows,
key=lambda row: (
row.rank_change(query.rank_change_days) is None,
-(row.rank_change(query.rank_change_days) or 0),
row.observation.sector_code,
),
)
)
else:
selected = select_rank_change_side(
metric_rows,
days=query.rank_change_days,
side=query.side,
)
else:
selected = select_percentile_side(metric_rows, query.side)
if query.side is not RankSide.BOTTOM:
selected = tuple(
sorted(
selected,
key=lambda row: (
row.rank_position is None,
row.rank_position or 0,
row.observation.sector_code,
),
)
)
search = query.search.strip().casefold() if query.search else ""
searched = tuple(
row
for row in selected
if not search
or search in row.observation.sector_code.casefold()
or search in row.observation.sector_name.casefold()
)
start = (query.page - 1) * query.page_size
return RankingPage(
status="success",
query=query,
publication=publication,
definition=definition,
rows=searched[start : start + query.page_size],
total=len(searched),
)
__all__ = [
"RadarDateIndex",
"RadarMetricDefinition",
"RadarQuery",
"RadarView",
"RankingPage",
"ReadSectorRadar",
]
@@ -226,8 +226,14 @@ class SectorRadarRepository(Protocol):
self, target_trade_date: date | None = None
) -> RadarPublication | None: ...
def get_successful_publication(self, target_trade_date: date) -> RadarPublication | None: ...
def get_latest_publication(self) -> RadarPublication | None: ...
def list_successful_dates(self) -> Sequence[date]: ...
def load_rankings(self, publication_id: str) -> Sequence[RankedMetric]: ...
def load_daily_aggregate_history(
self, target_trade_date: date, *, limit_dates: int
) -> Sequence[SectorDailyAggregate]: ...
@@ -279,7 +279,7 @@ class InMemorySectorRadarRepository:
) -> RadarPublication | None:
"""Find an identical success or partial revision without hiding failures."""
return next(
return max(
(
item
for item in self.publications.values()
@@ -287,7 +287,8 @@ class InMemorySectorRadarRepository:
and item.target_trade_date == target_trade_date
and item.input_hash == input_hash
),
None,
key=lambda item: (item.finished_at or item.started_at, item.publication_id),
default=None,
)
def get_last_good_publication(
@@ -303,7 +304,35 @@ class InMemorySectorRadarRepository:
)
return max(
candidates,
key=lambda item: (item.target_trade_date, item.finished_at or item.started_at),
key=lambda item: (
item.target_trade_date,
item.finished_at or item.started_at,
item.publication_id,
),
default=None,
)
def get_successful_publication(self, target_trade_date: date) -> RadarPublication | None:
"""Return the latest successful revision for exactly one date."""
candidates = tuple(
publication
for publication in self.publications.values()
if publication.status is PublicationStatus.SUCCESS
and publication.target_trade_date == target_trade_date
)
return max(
candidates,
key=lambda item: (item.finished_at or item.started_at, item.publication_id),
default=None,
)
def get_latest_publication(self) -> RadarPublication | None:
"""Return the newest build attempt regardless of terminal status."""
return max(
self.publications.values(),
key=lambda item: (item.target_trade_date, item.started_at, item.publication_id),
default=None,
)
@@ -321,6 +350,26 @@ class InMemorySectorRadarRepository:
)
)
def load_rankings(self, publication_id: str) -> Sequence[RankedMetric]:
"""Load every ranking projection owned by one publication."""
return tuple(
sorted(
(
record.ranking
for record in self.rankings.values()
if record.publication_id == publication_id
),
key=lambda row: (
row.observation.sector_type.value,
row.observation.metric_version,
row.rank_position is None,
row.rank_position or 0,
row.observation.sector_code,
),
)
)
def load_daily_aggregate_history(
self, target_trade_date: date, *, limit_dates: int
) -> Sequence[SectorDailyAggregate]:
@@ -378,7 +427,7 @@ class InMemorySectorRadarRepository:
for item in self.publications.values()
if item.status is PublicationStatus.SUCCESS and item.target_trade_date == trade_date
),
key=lambda item: item.finished_at or item.started_at,
key=lambda item: (item.finished_at or item.started_at, item.publication_id),
)
@staticmethod
@@ -645,7 +645,7 @@ class PostgresSectorRadarRepository:
row = connection.execute(
self._publication_select() + " WHERE target_trade_date = %s AND input_hash = %s "
"AND status IN ('success', 'partial') "
"ORDER BY finished_at DESC, created_at DESC, id DESC LIMIT 1",
"ORDER BY finished_at DESC, id DESC LIMIT 1",
(target_trade_date, input_hash),
).fetchone()
return None if row is None else self._publication_from_row(row)
@@ -663,12 +663,33 @@ class PostgresSectorRadarRepository:
query = (
self._publication_select()
+ where
+ " ORDER BY target_trade_date DESC, finished_at DESC, created_at DESC, id DESC LIMIT 1"
+ " ORDER BY target_trade_date DESC, finished_at DESC, id DESC LIMIT 1"
)
with self._connection() as connection:
row = connection.execute(query, parameters).fetchone()
return None if row is None else self._publication_from_row(row)
def get_successful_publication(self, target_trade_date: date) -> RadarPublication | None:
"""Read the latest successful revision for exactly one target date."""
with self._connection() as connection:
row = connection.execute(
self._publication_select() + " WHERE status = 'success' AND target_trade_date = %s "
"ORDER BY finished_at DESC, id DESC LIMIT 1",
(target_trade_date,),
).fetchone()
return None if row is None else self._publication_from_row(row)
def get_latest_publication(self) -> RadarPublication | None:
"""Read the newest build attempt regardless of terminal status."""
with self._connection() as connection:
row = connection.execute(
self._publication_select()
+ " ORDER BY target_trade_date DESC, started_at DESC, id DESC LIMIT 1"
).fetchone()
return None if row is None else self._publication_from_row(row)
def list_successful_dates(self) -> Sequence[date]:
"""List distinct successful target dates newest first."""
@@ -683,6 +704,24 @@ class PostgresSectorRadarRepository:
).fetchall()
return tuple(row[0] for row in rows)
def load_rankings(self, publication_id: str) -> Sequence[RankedMetric]:
"""Load all ranking projections for one publication in deterministic order."""
with self._connection() as connection:
rows = connection.execute(
"""
SELECT trade_date, sector_type, sector_code, sector_name, metric_kind,
metric_version, implementation_kind, unit, metric_value, quality,
member_count, valid_sample_count, membership_coverage,
moneyflow_coverage, rank_position, rank_percentile, rank_changes
FROM sector_radar_ranking
WHERE publication_id = %s
ORDER BY sector_type, metric_version, rank_position NULLS LAST, sector_code
""",
(publication_id,),
).fetchall()
return tuple(self._ranking_from_row(row) for row in rows)
def load_daily_aggregate_history(
self, target_trade_date: date, *, limit_dates: int
) -> Sequence[SectorDailyAggregate]:
@@ -697,7 +736,7 @@ class PostgresSectorRadarRepository:
SELECT DISTINCT ON (target_trade_date) id, target_trade_date
FROM sector_radar_publication
WHERE status = 'success' AND target_trade_date < %s
ORDER BY target_trade_date DESC, finished_at DESC, created_at DESC, id DESC
ORDER BY target_trade_date DESC, finished_at DESC, id DESC
LIMIT %s
)
SELECT aggregate.trade_date, aggregate.sector_type, aggregate.sector_code,
@@ -727,7 +766,7 @@ class PostgresSectorRadarRepository:
SELECT DISTINCT ON (target_trade_date) id, target_trade_date
FROM sector_radar_publication
WHERE status = 'success' AND target_trade_date < %s
ORDER BY target_trade_date DESC, finished_at DESC, created_at DESC, id DESC
ORDER BY target_trade_date DESC, finished_at DESC, id DESC
LIMIT %s
)
SELECT selected.target_trade_date, ranking.trade_date, ranking.sector_type,
@@ -0,0 +1,302 @@
"""HTTP presentation for persisted sector radar rankings."""
from __future__ import annotations
import atexit
import threading
from datetime import date, datetime
from decimal import Decimal
from typing import Annotated, Literal
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel, Field
from ....bootstrap.config import Settings, get_settings
from ..application.read import (
RadarDateIndex,
RadarMetricDefinition,
RadarQuery,
RadarView,
RankingPage,
ReadSectorRadar,
)
from ..domain.models import (
MetricKind,
MetricQuality,
MetricUnit,
PublicationStatus,
RadarPublication,
RankedMetric,
RankSide,
SectorType,
)
from ..infrastructure.postgres import (
PostgresSectorRadarRepository,
SectorRadarRepositoryError,
)
sector_radar_router = APIRouter()
_REPOSITORY_CACHE_LOCK = threading.Lock()
_REPOSITORY_CACHE: dict[tuple[str, int], PostgresSectorRadarRepository] = {}
class RadarPublicationResponse(BaseModel):
"""Safe publication provenance and quality metadata."""
publication_id: str
target_trade_date: date
status: PublicationStatus
source_version: str
universe_version: str
metric_versions: list[str]
input_hash: str | None
coverage: Decimal = Field(ge=0, le=1)
started_at: datetime
finished_at: datetime | None
error_summary: str | None
class RadarDatesResponse(BaseModel):
"""Successful dates plus newest attempt and strict last-good metadata."""
status: Literal["success", "no_data"]
available_dates: list[date]
current_attempt: RadarPublicationResponse | None
last_good: RadarPublicationResponse | None
class RadarMetricDefinitionResponse(BaseModel):
"""Version and labeling for one independent metric implementation."""
metric_kind: MetricKind
metric_version: str
label: str
unit: MetricUnit
implementation_kind: Literal["independent"]
disclaimer: str
class RadarRankingRowResponse(BaseModel):
"""One sector ranking row with explicit null and unit semantics."""
trade_date: date
sector_type: SectorType
sector_code: str
sector_name: str
metric_kind: MetricKind
metric_version: str
implementation_kind: Literal["independent"]
unit: MetricUnit
metric_value: Decimal | None
quality: MetricQuality
member_count: int = Field(ge=0)
valid_sample_count: int = Field(ge=0)
membership_coverage: Decimal = Field(ge=0, le=1)
moneyflow_coverage: Decimal = Field(ge=0, le=1)
rank_position: int | None = Field(default=None, ge=1)
rank_percentile: Decimal | None = Field(default=None, ge=0, le=100)
rank_change_days: int = Field(ge=1, le=5)
rank_change: int | None
def _empty_ranking_rows() -> list[RadarRankingRowResponse]:
return []
class RadarRankingsResponse(BaseModel):
"""One persisted, filtered ranking page."""
status: Literal["success", "no_data"]
requested_trade_date: date | None
sector_type: SectorType
view: RadarView
rank_change_metric: MetricKind
rank_change_days: int = Field(ge=1, le=5)
side: RankSide
search: str | None
publication: RadarPublicationResponse | None
definition: RadarMetricDefinitionResponse
page: int = Field(ge=1)
page_size: int = Field(ge=1, le=100)
total: int = Field(ge=0)
rows: list[RadarRankingRowResponse] = Field(default_factory=_empty_ranking_rows)
def get_sector_radar_reader(
settings: Annotated[Settings, Depends(get_settings)],
) -> ReadSectorRadar:
"""Return a reader backed by one process-cached PostgreSQL repository."""
key = (settings.database_url, settings.sector_radar_advisory_lock_key)
with _REPOSITORY_CACHE_LOCK:
repository = _REPOSITORY_CACHE.get(key)
if repository is None:
repository = PostgresSectorRadarRepository(
settings.database_url,
advisory_lock_key=settings.sector_radar_advisory_lock_key,
)
_REPOSITORY_CACHE[key] = repository
return ReadSectorRadar(repository)
def _close_cached_repositories() -> None:
"""Close process-owned radar pools during interpreter shutdown."""
with _REPOSITORY_CACHE_LOCK:
repositories = tuple(_REPOSITORY_CACHE.values())
_REPOSITORY_CACHE.clear()
for repository in repositories:
repository.close()
atexit.register(_close_cached_repositories)
@sector_radar_router.get("/dates", response_model=RadarDatesResponse)
def get_sector_radar_dates(
reader: Annotated[ReadSectorRadar, Depends(get_sector_radar_reader)],
) -> RadarDatesResponse:
"""Return persisted availability without invoking Tushare."""
try:
return _dates_response(reader.list_dates())
except SectorRadarRepositoryError as exc:
raise _storage_error() from exc
@sector_radar_router.get("/rankings", response_model=RadarRankingsResponse)
def get_sector_radar_rankings(
reader: Annotated[ReadSectorRadar, Depends(get_sector_radar_reader)],
trade_date: date | None = None,
sector_type: SectorType = SectorType.CONCEPT,
view: RadarView = RadarView.AMOUNT,
rank_change_metric: MetricKind = MetricKind.AMOUNT,
rank_change_days: Annotated[int, Query(ge=1, le=5)] = 1,
side: RankSide = RankSide.ALL,
search: Annotated[str | None, Query(max_length=100)] = None,
page: Annotated[int, Query(ge=1)] = 1,
page_size: Annotated[int, Query(ge=1, le=100)] = 20,
) -> RadarRankingsResponse:
"""Return one filtered page from a successful publication."""
query = RadarQuery(
trade_date=trade_date,
sector_type=sector_type,
view=view,
rank_change_metric=rank_change_metric,
rank_change_days=rank_change_days,
side=side,
search=search.strip() or None if search else None,
page=page,
page_size=page_size,
)
try:
return _rankings_response(reader.query(query))
except SectorRadarRepositoryError as exc:
raise _storage_error() from exc
def _dates_response(index: RadarDateIndex) -> RadarDatesResponse:
return RadarDatesResponse(
status=index.status,
available_dates=list(index.available_dates),
current_attempt=(
_publication_response(index.current_attempt)
if index.current_attempt is not None
else None
),
last_good=(_publication_response(index.last_good) if index.last_good is not None else None),
)
def _rankings_response(page: RankingPage) -> RadarRankingsResponse:
query = page.query
return RadarRankingsResponse(
status=page.status,
requested_trade_date=query.trade_date,
sector_type=query.sector_type,
view=query.view,
rank_change_metric=query.rank_change_metric,
rank_change_days=query.rank_change_days,
side=query.side,
search=query.search,
publication=(
_publication_response(page.publication) if page.publication is not None else None
),
definition=_definition_response(page.definition),
page=query.page,
page_size=query.page_size,
total=page.total,
rows=[_ranking_response(row, query.rank_change_days) for row in page.rows],
)
def _publication_response(publication: RadarPublication) -> RadarPublicationResponse:
return RadarPublicationResponse(
publication_id=publication.publication_id,
target_trade_date=publication.target_trade_date,
status=publication.status,
source_version=publication.source_version,
universe_version=publication.universe_version,
metric_versions=list(publication.metric_versions),
input_hash=publication.input_hash,
coverage=publication.coverage,
started_at=publication.started_at,
finished_at=publication.finished_at,
error_summary=publication.error_summary,
)
def _definition_response(
definition: RadarMetricDefinition,
) -> RadarMetricDefinitionResponse:
return RadarMetricDefinitionResponse(
metric_kind=definition.metric_kind,
metric_version=definition.metric_version,
label=definition.label,
unit=definition.unit,
implementation_kind=definition.implementation_kind,
disclaimer=definition.disclaimer,
)
def _ranking_response(row: RankedMetric, rank_change_days: int) -> RadarRankingRowResponse:
observation = row.observation
return RadarRankingRowResponse(
trade_date=observation.trade_date,
sector_type=observation.sector_type,
sector_code=observation.sector_code,
sector_name=observation.sector_name,
metric_kind=observation.metric_kind,
metric_version=observation.metric_version,
implementation_kind=observation.implementation_kind,
unit=observation.unit,
metric_value=observation.value,
quality=observation.quality,
member_count=observation.member_count,
valid_sample_count=observation.valid_sample_count,
membership_coverage=observation.membership_coverage,
moneyflow_coverage=observation.moneyflow_coverage,
rank_position=row.rank_position,
rank_percentile=row.rank_percentile,
rank_change_days=rank_change_days,
rank_change=row.rank_change(rank_change_days),
)
def _storage_error() -> HTTPException:
return HTTPException(
status_code=503,
detail={
"code": "sector_radar_storage_unavailable",
"message": "sector radar storage is unavailable",
},
)
__all__ = [
"RadarDatesResponse",
"RadarRankingsResponse",
"get_sector_radar_reader",
"sector_radar_router",
]