2026-09-05 10:22:08 +08:00
|
|
|
"""Application orchestration for persisted whole-universe strategy runs."""
|
2026-08-09 09:34:46 +08:00
|
|
|
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
import logging
|
2026-08-12 09:45:16 +08:00
|
|
|
import time
|
2026-09-05 10:22:08 +08:00
|
|
|
from collections import Counter
|
|
|
|
|
from collections.abc import Callable, Mapping, Sequence
|
2026-08-12 09:45:16 +08:00
|
|
|
from concurrent.futures import ThreadPoolExecutor
|
2026-09-05 19:48:30 +08:00
|
|
|
from dataclasses import dataclass, replace
|
2026-08-09 09:34:46 +08:00
|
|
|
from datetime import date
|
2026-09-05 10:22:08 +08:00
|
|
|
from typing import Protocol, cast
|
2026-08-09 09:34:46 +08:00
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
from ..domain.models import SelectionEvaluation, SelectionStrategyName, StockHistory
|
2026-08-31 16:14:16 +08:00
|
|
|
from ..domain.pattern_scoring import (
|
|
|
|
|
PatternCase,
|
|
|
|
|
PatternCaseLibraryLoader,
|
|
|
|
|
PatternScore,
|
|
|
|
|
PatternScorer,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
from ..domain.runs import (
|
2026-08-12 09:45:16 +08:00
|
|
|
BatchSelectionRunStore,
|
|
|
|
|
BatchSelectionUniverseReader,
|
2026-08-09 09:34:46 +08:00
|
|
|
SelectionExecutionSource,
|
|
|
|
|
SelectionRerunRequired,
|
2026-08-10 11:09:23 +08:00
|
|
|
SelectionResultQuery,
|
2026-08-09 09:34:46 +08:00
|
|
|
SelectionRun,
|
|
|
|
|
SelectionRunInProgress,
|
|
|
|
|
SelectionRunItem,
|
|
|
|
|
SelectionRunStatus,
|
|
|
|
|
SelectionRunStore,
|
2026-09-05 19:48:30 +08:00
|
|
|
SelectionSectorCount,
|
|
|
|
|
SelectionSectorMembership,
|
|
|
|
|
SelectionSectorReader,
|
2026-08-12 09:45:16 +08:00
|
|
|
SelectionStock,
|
2026-08-09 09:34:46 +08:00
|
|
|
SelectionUniverseReader,
|
|
|
|
|
)
|
|
|
|
|
from .evaluate import EvaluateZhixingB1
|
|
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
# Selection runs are started by the ASGI service in production. A child of
|
|
|
|
|
# Uvicorn's configured logger keeps INFO diagnostics visible in container logs.
|
|
|
|
|
logger = logging.getLogger("uvicorn.error.zhixing.selection.run")
|
|
|
|
|
StrategyName = SelectionStrategyName
|
|
|
|
|
_FAILURE_STATUSES = {
|
|
|
|
|
"insufficient_history",
|
|
|
|
|
"missing_target_bar",
|
|
|
|
|
"missing_turnover_rate",
|
|
|
|
|
"data_error",
|
|
|
|
|
}
|
2026-08-09 09:34:46 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
class SelectionEvaluator(Protocol):
|
|
|
|
|
"""Minimal single-stock evaluator required by the batch orchestrator."""
|
|
|
|
|
|
|
|
|
|
def execute(self, ts_code: str, target_trade_date: date) -> SelectionEvaluation: ...
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
|
|
|
class PreparedSelectionRun:
|
|
|
|
|
"""A claimed run and its immutable market-data source snapshot."""
|
|
|
|
|
|
|
|
|
|
run: SelectionRun
|
|
|
|
|
source: SelectionExecutionSource
|
|
|
|
|
|
|
|
|
|
|
2026-09-05 19:48:30 +08:00
|
|
|
@dataclass(frozen=True, slots=True)
|
|
|
|
|
class SelectionSectorAggregates:
|
|
|
|
|
"""A run summary plus its selected stocks' sector membership counts."""
|
|
|
|
|
|
|
|
|
|
run: SelectionRun | None
|
|
|
|
|
snapshot_trade_date: date | None
|
|
|
|
|
sector_type: str
|
|
|
|
|
sectors: tuple[SelectionSectorCount, ...]
|
|
|
|
|
|
|
|
|
|
|
2026-08-09 09:34:46 +08:00
|
|
|
class RunZhixingB1:
|
2026-09-05 10:22:08 +08:00
|
|
|
"""Prepare, execute, and query persisted selection strategy batches.
|
|
|
|
|
|
|
|
|
|
The historical class name remains as a compatibility seam for existing
|
|
|
|
|
composition and tests while strategy routing is now explicit.
|
|
|
|
|
"""
|
2026-08-09 09:34:46 +08:00
|
|
|
|
|
|
|
|
def __init__(
|
|
|
|
|
self,
|
|
|
|
|
reader: SelectionUniverseReader,
|
|
|
|
|
store: SelectionRunStore,
|
|
|
|
|
evaluator: SelectionEvaluator | None = None,
|
2026-09-05 10:22:08 +08:00
|
|
|
evaluators: Mapping[StrategyName, SelectionEvaluator] | None = None,
|
2026-08-31 16:14:16 +08:00
|
|
|
pattern_case_loader: PatternCaseLibraryLoader | None = None,
|
|
|
|
|
pattern_scorer: PatternScorer | None = None,
|
2026-09-05 19:48:30 +08:00
|
|
|
sector_reader: SelectionSectorReader | None = None,
|
2026-08-12 09:45:16 +08:00
|
|
|
*,
|
2026-08-31 16:14:16 +08:00
|
|
|
pattern_scoring_enabled: bool = False,
|
2026-08-12 09:45:16 +08:00
|
|
|
max_workers: int = 4,
|
|
|
|
|
batch_size: int = 200,
|
2026-08-09 09:34:46 +08:00
|
|
|
) -> None:
|
2026-08-12 09:45:16 +08:00
|
|
|
"""Inject storage ports and configure bounded chunk execution."""
|
2026-08-09 09:34:46 +08:00
|
|
|
|
2026-08-12 09:45:16 +08:00
|
|
|
if max_workers < 1:
|
|
|
|
|
raise ValueError("max_workers must be at least 1")
|
|
|
|
|
if batch_size < 1:
|
|
|
|
|
raise ValueError("batch_size must be at least 1")
|
2026-08-09 09:34:46 +08:00
|
|
|
self.reader = reader
|
|
|
|
|
self.store = store
|
|
|
|
|
self.evaluator = evaluator or EvaluateZhixingB1(reader)
|
2026-09-05 10:22:08 +08:00
|
|
|
self.evaluators: dict[StrategyName, SelectionEvaluator] = {
|
|
|
|
|
"zhixing_b1": self.evaluator,
|
|
|
|
|
}
|
|
|
|
|
if evaluators is not None:
|
|
|
|
|
self.evaluators.update(evaluators)
|
2026-08-31 16:14:16 +08:00
|
|
|
self.pattern_case_loader = pattern_case_loader
|
|
|
|
|
self.pattern_scorer = pattern_scorer
|
2026-09-05 19:48:30 +08:00
|
|
|
self.sector_reader = sector_reader
|
2026-08-31 16:14:16 +08:00
|
|
|
self.pattern_scoring_enabled = pattern_scoring_enabled
|
2026-08-12 09:45:16 +08:00
|
|
|
self.max_workers = max_workers
|
|
|
|
|
self.batch_size = batch_size
|
2026-08-09 09:34:46 +08:00
|
|
|
|
|
|
|
|
def prepare(
|
|
|
|
|
self,
|
|
|
|
|
strategy: StrategyName,
|
|
|
|
|
target_trade_date: date,
|
|
|
|
|
*,
|
|
|
|
|
rerun: bool,
|
|
|
|
|
) -> PreparedSelectionRun:
|
|
|
|
|
"""Validate source eligibility before claiming the rerunnable key."""
|
|
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
logger.info(
|
|
|
|
|
"selection_run_prepare_started strategy=%s target_trade_date=%s rerun=%s",
|
2026-08-09 09:34:46 +08:00
|
|
|
strategy,
|
2026-09-05 10:22:08 +08:00
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
rerun,
|
|
|
|
|
)
|
|
|
|
|
if strategy not in self.evaluators:
|
|
|
|
|
raise ValueError(f"selection evaluator is not configured for strategy: {strategy}")
|
|
|
|
|
try:
|
|
|
|
|
source = self.reader.load_execution_source(strategy, target_trade_date)
|
|
|
|
|
run = self.store.prepare_run(
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date,
|
|
|
|
|
source,
|
|
|
|
|
rerun=rerun,
|
|
|
|
|
)
|
|
|
|
|
except Exception as exc: # noqa: BLE001 - log the safe prepare boundary and preserve type
|
|
|
|
|
logger.warning(
|
|
|
|
|
"selection_run_prepare_failed strategy=%s target_trade_date=%s "
|
|
|
|
|
"status=failed error_type=%s reason=%s",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
exc.__class__.__name__,
|
|
|
|
|
_safe_item_error(exc),
|
|
|
|
|
)
|
|
|
|
|
raise
|
|
|
|
|
logger.info(
|
|
|
|
|
"selection_run_prepared strategy=%s target_trade_date=%s run_id=%s "
|
|
|
|
|
"market_sync_batch_id=%s target_count=%d eligible_count=%d "
|
|
|
|
|
"coverage=%s status=running",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
run.id,
|
|
|
|
|
source.market_sync_batch_id,
|
|
|
|
|
source.target_count,
|
|
|
|
|
len(source.stocks),
|
|
|
|
|
source.coverage,
|
2026-08-09 09:34:46 +08:00
|
|
|
)
|
|
|
|
|
return PreparedSelectionRun(run=run, source=source)
|
|
|
|
|
|
|
|
|
|
def execute(self, prepared: PreparedSelectionRun) -> None:
|
|
|
|
|
"""Evaluate every eligible stock and converge the persisted run status.
|
|
|
|
|
|
|
|
|
|
This method is the boundary used by FastAPI's in-process background
|
|
|
|
|
task. An unexpected batch-level error is recorded before the worker
|
|
|
|
|
returns so the UI never mistakes a lost worker exception for success.
|
|
|
|
|
"""
|
|
|
|
|
|
2026-08-12 09:45:16 +08:00
|
|
|
stocks = _unique_stocks(prepared.source.stocks)
|
2026-09-05 10:22:08 +08:00
|
|
|
strategy = prepared.run.strategy
|
|
|
|
|
target_trade_date = prepared.source.target_trade_date
|
|
|
|
|
evaluator = self.evaluators[strategy]
|
2026-08-09 09:34:46 +08:00
|
|
|
evaluated_count = 0
|
|
|
|
|
selected_stock_count = 0
|
|
|
|
|
signal_count = 0
|
|
|
|
|
failed_count = 0
|
2026-09-05 10:22:08 +08:00
|
|
|
missing_turnover_count = 0
|
|
|
|
|
insufficient_history_count = 0
|
2026-08-12 09:45:16 +08:00
|
|
|
history_rows = 0
|
|
|
|
|
batch_count = _chunk_count(len(stocks), self.batch_size)
|
2026-09-05 10:22:08 +08:00
|
|
|
current_batch = 0
|
|
|
|
|
final_status: SelectionRunStatus = "failed"
|
2026-08-12 09:45:16 +08:00
|
|
|
read_seconds = 0.0
|
|
|
|
|
evaluate_seconds = 0.0
|
|
|
|
|
persist_seconds = 0.0
|
2026-08-31 16:14:16 +08:00
|
|
|
scoring_seconds = 0.0
|
2026-09-05 10:22:08 +08:00
|
|
|
logger.info(
|
|
|
|
|
"selection_run_started strategy=%s target_trade_date=%s run_id=%s "
|
|
|
|
|
"market_sync_batch_id=%s stock_count=%d batch_count=%d worker_count=%d "
|
|
|
|
|
"status=running",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
prepared.source.market_sync_batch_id,
|
|
|
|
|
len(stocks),
|
|
|
|
|
batch_count,
|
|
|
|
|
self.max_workers,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
try:
|
2026-09-05 10:22:08 +08:00
|
|
|
if strategy == "zhixing_b1":
|
|
|
|
|
pattern_cases, pattern_library_error = self._prepare_pattern_cases(prepared.run.id)
|
|
|
|
|
else:
|
|
|
|
|
pattern_cases, pattern_library_error = None, None
|
2026-08-12 09:45:16 +08:00
|
|
|
with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
|
2026-09-05 10:22:08 +08:00
|
|
|
for batch_index, batch_stocks in enumerate(
|
|
|
|
|
_chunks(stocks, self.batch_size),
|
|
|
|
|
start=1,
|
|
|
|
|
):
|
|
|
|
|
current_batch = batch_index
|
2026-08-12 09:45:16 +08:00
|
|
|
read_started = time.perf_counter()
|
|
|
|
|
histories = self._load_histories(
|
|
|
|
|
batch_stocks,
|
2026-09-05 10:22:08 +08:00
|
|
|
target_trade_date,
|
|
|
|
|
evaluator,
|
2026-08-09 09:34:46 +08:00
|
|
|
)
|
2026-09-05 10:22:08 +08:00
|
|
|
batch_read_seconds = time.perf_counter() - read_started
|
|
|
|
|
read_seconds += batch_read_seconds
|
|
|
|
|
batch_history_rows = sum(
|
2026-08-12 09:45:16 +08:00
|
|
|
len(history.bars) for history in histories if history is not None
|
2026-08-09 09:34:46 +08:00
|
|
|
)
|
2026-09-05 10:22:08 +08:00
|
|
|
history_rows += batch_history_rows
|
|
|
|
|
batch_missing_turnover = sum(
|
2026-09-05 19:48:30 +08:00
|
|
|
not _turnover_present(history, target_trade_date) for history in histories
|
2026-09-05 10:22:08 +08:00
|
|
|
)
|
|
|
|
|
logger.info(
|
|
|
|
|
"selection_read_batch_summary strategy=%s target_trade_date=%s "
|
|
|
|
|
"run_id=%s batch=%d batch_count=%d stock_count=%d history_rows=%d "
|
|
|
|
|
"turnover_missing_count=%d status=success read_seconds=%.3f",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
batch_index,
|
|
|
|
|
batch_count,
|
|
|
|
|
len(batch_stocks),
|
|
|
|
|
batch_history_rows,
|
|
|
|
|
batch_missing_turnover,
|
|
|
|
|
batch_read_seconds,
|
|
|
|
|
)
|
2026-08-12 09:45:16 +08:00
|
|
|
|
|
|
|
|
evaluate_started = time.perf_counter()
|
2026-08-31 16:14:16 +08:00
|
|
|
evaluations = tuple(
|
|
|
|
|
executor.map(
|
|
|
|
|
self._evaluate_stock,
|
|
|
|
|
batch_stocks,
|
|
|
|
|
histories,
|
2026-09-05 10:22:08 +08:00
|
|
|
[target_trade_date] * len(batch_stocks),
|
|
|
|
|
[evaluator] * len(batch_stocks),
|
|
|
|
|
[strategy] * len(batch_stocks),
|
|
|
|
|
[prepared.run.id] * len(batch_stocks),
|
|
|
|
|
[batch_index] * len(batch_stocks),
|
2026-08-31 16:14:16 +08:00
|
|
|
)
|
|
|
|
|
)
|
2026-09-05 10:22:08 +08:00
|
|
|
batch_evaluate_seconds = time.perf_counter() - evaluate_started
|
|
|
|
|
evaluate_seconds += batch_evaluate_seconds
|
|
|
|
|
status_counts = Counter(evaluation.status for evaluation in evaluations)
|
|
|
|
|
no_signal_reasons = Counter(
|
|
|
|
|
evaluation.reason or "unspecified"
|
|
|
|
|
for evaluation in evaluations
|
|
|
|
|
if evaluation.status == "no_signal"
|
|
|
|
|
)
|
|
|
|
|
missing_turnover_count += status_counts["missing_turnover_rate"]
|
|
|
|
|
insufficient_history_count += status_counts["insufficient_history"]
|
|
|
|
|
for stock, history, evaluation in zip(
|
|
|
|
|
batch_stocks,
|
|
|
|
|
histories,
|
|
|
|
|
evaluations,
|
|
|
|
|
strict=True,
|
|
|
|
|
):
|
|
|
|
|
if evaluation.status not in _FAILURE_STATUSES:
|
|
|
|
|
continue
|
|
|
|
|
logger.warning(
|
|
|
|
|
"selection_item_incomplete strategy=%s target_trade_date=%s "
|
|
|
|
|
"run_id=%s batch=%d ts_code=%s history_rows=%d "
|
|
|
|
|
"turnover_present=%s status=%s error_type=%s reason=%s",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
batch_index,
|
|
|
|
|
stock.ts_code,
|
|
|
|
|
len(history.bars) if history is not None else 0,
|
|
|
|
|
_turnover_present(history, target_trade_date),
|
|
|
|
|
evaluation.status,
|
|
|
|
|
evaluation.status,
|
|
|
|
|
evaluation.reason or evaluation.status,
|
|
|
|
|
)
|
2026-08-31 16:14:16 +08:00
|
|
|
|
|
|
|
|
scoring_started = time.perf_counter()
|
2026-08-12 09:45:16 +08:00
|
|
|
items = tuple(
|
|
|
|
|
_to_item(
|
|
|
|
|
stock.ts_code,
|
|
|
|
|
stock.name,
|
|
|
|
|
evaluation,
|
2026-08-31 16:14:16 +08:00
|
|
|
pattern_score=self._score_stock(
|
2026-09-05 10:22:08 +08:00
|
|
|
strategy,
|
2026-08-31 16:14:16 +08:00
|
|
|
prepared.run.id,
|
|
|
|
|
stock,
|
|
|
|
|
history,
|
|
|
|
|
evaluation,
|
|
|
|
|
pattern_cases,
|
|
|
|
|
pattern_library_error,
|
|
|
|
|
),
|
2026-08-12 09:45:16 +08:00
|
|
|
)
|
2026-08-31 16:14:16 +08:00
|
|
|
for stock, history, evaluation in zip(
|
2026-08-12 09:45:16 +08:00
|
|
|
batch_stocks,
|
2026-08-31 16:14:16 +08:00
|
|
|
histories,
|
|
|
|
|
evaluations,
|
2026-08-12 09:45:16 +08:00
|
|
|
strict=True,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
)
|
2026-08-31 16:14:16 +08:00
|
|
|
scoring_seconds += time.perf_counter() - scoring_started
|
2026-08-12 09:45:16 +08:00
|
|
|
|
|
|
|
|
evaluated_count += len(items)
|
|
|
|
|
selected_stock_count += sum(item.status == "selected" for item in items)
|
|
|
|
|
signal_count += sum(item.signal_count for item in items)
|
|
|
|
|
failed_count += sum(item.status in _FAILURE_STATUSES for item in items)
|
|
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
logger.info(
|
|
|
|
|
"selection_evaluate_batch_summary strategy=%s target_trade_date=%s "
|
|
|
|
|
"run_id=%s batch=%d batch_count=%d stock_count=%d selected_count=%d "
|
|
|
|
|
"no_signal_count=%d insufficient_history_count=%d "
|
|
|
|
|
"missing_target_bar_count=%d missing_turnover_count=%d "
|
|
|
|
|
"data_error_count=%d no_signal_reasons=%s status=complete "
|
|
|
|
|
"evaluate_seconds=%.3f",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
batch_index,
|
|
|
|
|
batch_count,
|
|
|
|
|
len(items),
|
|
|
|
|
status_counts["selected"],
|
|
|
|
|
status_counts["no_signal"],
|
|
|
|
|
status_counts["insufficient_history"],
|
|
|
|
|
status_counts["missing_target_bar"],
|
|
|
|
|
status_counts["missing_turnover_rate"],
|
|
|
|
|
status_counts["data_error"],
|
|
|
|
|
dict(no_signal_reasons),
|
|
|
|
|
batch_evaluate_seconds,
|
|
|
|
|
)
|
|
|
|
|
|
2026-08-12 09:45:16 +08:00
|
|
|
persist_started = time.perf_counter()
|
|
|
|
|
self._record_items(prepared.run.id, items)
|
2026-09-05 10:22:08 +08:00
|
|
|
batch_persist_seconds = time.perf_counter() - persist_started
|
|
|
|
|
persist_seconds += batch_persist_seconds
|
|
|
|
|
logger.info(
|
|
|
|
|
"selection_persist_batch_summary strategy=%s target_trade_date=%s "
|
|
|
|
|
"run_id=%s batch=%d batch_count=%d item_count=%d signal_count=%d "
|
|
|
|
|
"status=success persist_seconds=%.3f",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
batch_index,
|
|
|
|
|
batch_count,
|
|
|
|
|
len(items),
|
|
|
|
|
sum(item.signal_count for item in items),
|
|
|
|
|
batch_persist_seconds,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
final_status = _run_status(evaluated_count, failed_count)
|
2026-08-09 09:34:46 +08:00
|
|
|
self.store.finish_run(
|
|
|
|
|
prepared.run.id,
|
2026-09-05 10:22:08 +08:00
|
|
|
final_status,
|
2026-08-09 09:34:46 +08:00
|
|
|
evaluated_count=evaluated_count,
|
|
|
|
|
selected_stock_count=selected_stock_count,
|
|
|
|
|
signal_count=signal_count,
|
|
|
|
|
failed_count=failed_count,
|
|
|
|
|
)
|
2026-09-05 10:22:08 +08:00
|
|
|
logger.info(
|
|
|
|
|
"selection_run_converged strategy=%s target_trade_date=%s run_id=%s "
|
|
|
|
|
"batch=%d evaluated_count=%d selected_stock_count=%d signal_count=%d "
|
|
|
|
|
"failed_count=%d status=%s",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
current_batch,
|
|
|
|
|
evaluated_count,
|
|
|
|
|
selected_stock_count,
|
|
|
|
|
signal_count,
|
|
|
|
|
failed_count,
|
|
|
|
|
final_status,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
except Exception as exc: # noqa: BLE001 - worker boundary must persist failure state
|
2026-09-05 10:22:08 +08:00
|
|
|
final_status = "failed"
|
2026-08-12 09:45:16 +08:00
|
|
|
logger.error(
|
2026-09-05 10:22:08 +08:00
|
|
|
"selection_run_failed strategy=%s target_trade_date=%s run_id=%s "
|
|
|
|
|
"batch=%d status=failed error_type=%s reason=%s",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
2026-08-12 09:45:16 +08:00
|
|
|
prepared.run.id,
|
2026-09-05 10:22:08 +08:00
|
|
|
current_batch,
|
2026-08-12 09:45:16 +08:00
|
|
|
exc.__class__.__name__,
|
|
|
|
|
_safe_item_error(exc),
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
try:
|
|
|
|
|
self.store.finish_run(
|
|
|
|
|
prepared.run.id,
|
|
|
|
|
"failed",
|
|
|
|
|
evaluated_count=evaluated_count,
|
|
|
|
|
selected_stock_count=selected_stock_count,
|
|
|
|
|
signal_count=signal_count,
|
|
|
|
|
failed_count=max(failed_count, 1),
|
|
|
|
|
error_type="batch_error",
|
2026-09-05 10:22:08 +08:00
|
|
|
error_message=_safe_item_error(exc),
|
2026-08-09 09:34:46 +08:00
|
|
|
)
|
|
|
|
|
except Exception: # noqa: BLE001 - preserve the original worker failure
|
2026-08-12 09:45:16 +08:00
|
|
|
logger.error(
|
2026-09-05 10:22:08 +08:00
|
|
|
"selection_run_failure_persist_failed strategy=%s target_trade_date=%s "
|
|
|
|
|
"run_id=%s batch=%d status=failed error_type=finish_run_failed",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
2026-08-12 09:45:16 +08:00
|
|
|
prepared.run.id,
|
2026-09-05 10:22:08 +08:00
|
|
|
current_batch,
|
2026-08-12 09:45:16 +08:00
|
|
|
)
|
|
|
|
|
finally:
|
|
|
|
|
logger.info(
|
2026-09-05 10:22:08 +08:00
|
|
|
"selection_run_summary strategy=%s target_trade_date=%s run_id=%s "
|
|
|
|
|
"market_sync_batch_id=%s stock_count=%d history_rows=%d batch_count=%d "
|
|
|
|
|
"last_batch=%d worker_count=%d evaluated_count=%d selected_stock_count=%d "
|
|
|
|
|
"signal_count=%d failed_count=%d insufficient_history_count=%d "
|
|
|
|
|
"missing_turnover_count=%d status=%s read_seconds=%.3f "
|
2026-08-31 16:14:16 +08:00
|
|
|
"evaluate_seconds=%.3f scoring_seconds=%.3f persist_seconds=%.3f",
|
2026-09-05 10:22:08 +08:00
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
2026-08-12 09:45:16 +08:00
|
|
|
prepared.run.id,
|
2026-09-05 10:22:08 +08:00
|
|
|
prepared.source.market_sync_batch_id,
|
2026-08-12 09:45:16 +08:00
|
|
|
len(stocks),
|
|
|
|
|
history_rows,
|
|
|
|
|
batch_count,
|
2026-09-05 10:22:08 +08:00
|
|
|
current_batch,
|
2026-08-12 09:45:16 +08:00
|
|
|
self.max_workers,
|
2026-09-05 10:22:08 +08:00
|
|
|
evaluated_count,
|
|
|
|
|
selected_stock_count,
|
|
|
|
|
signal_count,
|
|
|
|
|
failed_count,
|
|
|
|
|
insufficient_history_count,
|
|
|
|
|
missing_turnover_count,
|
|
|
|
|
final_status,
|
2026-08-12 09:45:16 +08:00
|
|
|
read_seconds,
|
|
|
|
|
evaluate_seconds,
|
2026-08-31 16:14:16 +08:00
|
|
|
scoring_seconds,
|
2026-08-12 09:45:16 +08:00
|
|
|
persist_seconds,
|
|
|
|
|
)
|
|
|
|
|
|
2026-08-31 16:14:16 +08:00
|
|
|
def _prepare_pattern_cases(
|
|
|
|
|
self,
|
|
|
|
|
run_id: str,
|
|
|
|
|
) -> tuple[tuple[PatternCase, ...] | None, str | None]:
|
|
|
|
|
"""Load the complete case library once without failing selection."""
|
|
|
|
|
|
|
|
|
|
if not self.pattern_scoring_enabled:
|
|
|
|
|
return None, None
|
|
|
|
|
if self.pattern_case_loader is None or self.pattern_scorer is None:
|
|
|
|
|
reason = "pattern scoring is enabled but not configured"
|
|
|
|
|
logger.error("selection_pattern_library_failed run_id=%s reason=%s", run_id, reason)
|
|
|
|
|
return None, reason
|
|
|
|
|
try:
|
|
|
|
|
return self.pattern_case_loader.load(), None
|
|
|
|
|
except Exception as exc: # noqa: BLE001 - scoring enrichment must not fail selection
|
|
|
|
|
reason = _safe_item_error(exc)
|
|
|
|
|
logger.warning(
|
|
|
|
|
"selection_pattern_library_failed run_id=%s error_type=%s reason=%s",
|
|
|
|
|
run_id,
|
|
|
|
|
exc.__class__.__name__,
|
|
|
|
|
reason,
|
|
|
|
|
)
|
|
|
|
|
return None, reason
|
|
|
|
|
|
|
|
|
|
def _score_stock(
|
|
|
|
|
self,
|
2026-09-05 10:22:08 +08:00
|
|
|
strategy: StrategyName,
|
2026-08-31 16:14:16 +08:00
|
|
|
run_id: str,
|
|
|
|
|
stock: SelectionStock,
|
|
|
|
|
history: StockHistory | None,
|
|
|
|
|
evaluation: SelectionEvaluation,
|
|
|
|
|
cases: tuple[PatternCase, ...] | None,
|
|
|
|
|
library_error: str | None,
|
|
|
|
|
) -> PatternScore:
|
|
|
|
|
"""Score one selected stock once and isolate enrichment failures."""
|
|
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
if (
|
|
|
|
|
strategy != "zhixing_b1"
|
|
|
|
|
or not self.pattern_scoring_enabled
|
|
|
|
|
or evaluation.status != "selected"
|
|
|
|
|
):
|
2026-08-31 16:14:16 +08:00
|
|
|
return PatternScore()
|
|
|
|
|
if library_error is not None:
|
|
|
|
|
return PatternScore.failed(library_error)
|
|
|
|
|
if history is None or cases is None or self.pattern_scorer is None:
|
|
|
|
|
return PatternScore.failed("pattern scoring history or case library is unavailable")
|
|
|
|
|
try:
|
|
|
|
|
return self.pattern_scorer.score(history, cases)
|
|
|
|
|
except Exception as exc: # noqa: BLE001 - one score must not fail the selection run
|
|
|
|
|
reason = _safe_item_error(exc)
|
|
|
|
|
logger.warning(
|
|
|
|
|
"selection_pattern_score_failed run_id=%s ts_code=%s error_type=%s reason=%s",
|
|
|
|
|
run_id,
|
|
|
|
|
stock.ts_code,
|
|
|
|
|
exc.__class__.__name__,
|
|
|
|
|
reason,
|
|
|
|
|
)
|
|
|
|
|
return PatternScore.failed(reason)
|
|
|
|
|
|
2026-08-12 09:45:16 +08:00
|
|
|
def _load_histories(
|
|
|
|
|
self,
|
|
|
|
|
stocks: Sequence[SelectionStock],
|
|
|
|
|
target_trade_date: date,
|
2026-09-05 10:22:08 +08:00
|
|
|
evaluator: SelectionEvaluator,
|
2026-08-12 09:45:16 +08:00
|
|
|
) -> tuple[StockHistory | None, ...]:
|
|
|
|
|
"""Load one chunk when the reader supports it, with old-path fallback."""
|
|
|
|
|
|
|
|
|
|
typed_stocks = tuple(stocks)
|
|
|
|
|
loader = getattr(self.reader, "load_histories", None)
|
|
|
|
|
if callable(loader):
|
|
|
|
|
batch_reader = cast(BatchSelectionUniverseReader, self.reader)
|
|
|
|
|
loaded = batch_reader.load_histories(typed_stocks, target_trade_date)
|
|
|
|
|
histories_by_code = {history.ts_code: history for history in loaded}
|
|
|
|
|
return tuple(
|
|
|
|
|
histories_by_code.get(
|
|
|
|
|
stock.ts_code,
|
|
|
|
|
StockHistory(ts_code=stock.ts_code, name=stock.name),
|
|
|
|
|
)
|
|
|
|
|
for stock in typed_stocks
|
|
|
|
|
)
|
|
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
execute_history = getattr(evaluator, "execute_history", None)
|
|
|
|
|
if callable(execute_history):
|
2026-08-12 09:45:16 +08:00
|
|
|
return tuple(
|
|
|
|
|
self.reader.load_history(stock.ts_code, target_trade_date) for stock in typed_stocks
|
|
|
|
|
)
|
|
|
|
|
return (None,) * len(typed_stocks)
|
|
|
|
|
|
|
|
|
|
def _evaluate_stock(
|
|
|
|
|
self,
|
|
|
|
|
stock: SelectionStock,
|
|
|
|
|
history: StockHistory | None,
|
|
|
|
|
target_trade_date: date,
|
2026-09-05 10:22:08 +08:00
|
|
|
evaluator: SelectionEvaluator,
|
|
|
|
|
strategy: StrategyName,
|
|
|
|
|
run_id: str,
|
|
|
|
|
batch_index: int,
|
2026-08-12 09:45:16 +08:00
|
|
|
) -> SelectionEvaluation:
|
|
|
|
|
"""Evaluate one stock inside a worker and isolate its exception."""
|
|
|
|
|
|
|
|
|
|
ts_code = stock.ts_code
|
|
|
|
|
try:
|
|
|
|
|
execute_history: Callable[[StockHistory, date], SelectionEvaluation] | None = getattr(
|
2026-09-05 10:22:08 +08:00
|
|
|
evaluator,
|
2026-08-12 09:45:16 +08:00
|
|
|
"execute_history",
|
|
|
|
|
None,
|
|
|
|
|
)
|
|
|
|
|
if history is not None and execute_history is not None:
|
|
|
|
|
return execute_history(history, target_trade_date)
|
2026-09-05 10:22:08 +08:00
|
|
|
return evaluator.execute(ts_code, target_trade_date)
|
2026-08-12 09:45:16 +08:00
|
|
|
except Exception as exc: # noqa: BLE001 - isolate one stock from the batch
|
|
|
|
|
logger.warning(
|
2026-09-05 10:22:08 +08:00
|
|
|
"selection_item_failed strategy=%s target_trade_date=%s run_id=%s "
|
|
|
|
|
"batch=%d ts_code=%s history_rows=%d turnover_present=%s "
|
|
|
|
|
"status=data_error error_type=%s reason=%s",
|
|
|
|
|
strategy,
|
|
|
|
|
target_trade_date.isoformat(),
|
|
|
|
|
run_id,
|
|
|
|
|
batch_index,
|
2026-08-12 09:45:16 +08:00
|
|
|
ts_code,
|
2026-09-05 10:22:08 +08:00
|
|
|
len(history.bars) if history is not None else 0,
|
|
|
|
|
_turnover_present(history, target_trade_date),
|
2026-08-12 09:45:16 +08:00
|
|
|
exc.__class__.__name__,
|
|
|
|
|
_safe_item_error(exc),
|
|
|
|
|
)
|
|
|
|
|
return SelectionEvaluation(
|
|
|
|
|
ts_code=ts_code,
|
|
|
|
|
target_trade_date=target_trade_date,
|
|
|
|
|
status="data_error",
|
|
|
|
|
reason=_safe_item_error(exc),
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
def _record_items(self, run_id: str, items: Sequence[SelectionRunItem]) -> None:
|
|
|
|
|
"""Use batch persistence while retaining the old single-item seam."""
|
|
|
|
|
|
|
|
|
|
record_items = getattr(self.store, "record_items", None)
|
|
|
|
|
if callable(record_items):
|
|
|
|
|
batch_store = cast(BatchSelectionRunStore, self.store)
|
|
|
|
|
batch_store.record_items(run_id, tuple(items))
|
|
|
|
|
return
|
|
|
|
|
for item in items:
|
|
|
|
|
self.store.record_item(run_id, item)
|
2026-08-09 09:34:46 +08:00
|
|
|
|
2026-08-10 11:09:23 +08:00
|
|
|
def get_run(
|
|
|
|
|
self,
|
|
|
|
|
run_id: str,
|
|
|
|
|
*,
|
|
|
|
|
query: SelectionResultQuery | None = None,
|
|
|
|
|
) -> SelectionRun | None:
|
2026-09-05 19:48:30 +08:00
|
|
|
"""Read one persisted run for polling with optional sector filtering."""
|
|
|
|
|
|
|
|
|
|
effective = query or SelectionResultQuery()
|
|
|
|
|
if not effective.sector:
|
|
|
|
|
return self.store.get_run(run_id, query=query)
|
|
|
|
|
identity = self.store.get_run_identity(run_id)
|
|
|
|
|
if identity is None:
|
|
|
|
|
return None
|
|
|
|
|
member_codes = self._sector_member_codes(identity.target_trade_date, effective.sector)
|
|
|
|
|
if member_codes is None:
|
|
|
|
|
return self.store.get_run(run_id, query=query)
|
|
|
|
|
return self.store.get_run(
|
|
|
|
|
run_id,
|
|
|
|
|
query=replace(effective, sector=None),
|
|
|
|
|
sector_stock_codes=member_codes,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
|
|
|
|
|
def get_latest(
|
|
|
|
|
self,
|
|
|
|
|
strategy: StrategyName,
|
|
|
|
|
target_trade_date: date | None = None,
|
2026-08-10 11:09:23 +08:00
|
|
|
*,
|
|
|
|
|
query: SelectionResultQuery | None = None,
|
2026-08-09 09:34:46 +08:00
|
|
|
) -> SelectionRun | None:
|
|
|
|
|
"""Read the current result by date or the latest result for a strategy."""
|
|
|
|
|
|
2026-09-05 19:48:30 +08:00
|
|
|
effective = query or SelectionResultQuery()
|
|
|
|
|
if not effective.sector:
|
|
|
|
|
return self.store.get_latest_run(strategy, target_trade_date, query=query)
|
|
|
|
|
identity = self.store.get_latest_run_identity(strategy, target_trade_date)
|
|
|
|
|
if identity is None:
|
|
|
|
|
return None
|
|
|
|
|
member_codes = self._sector_member_codes(identity.target_trade_date, effective.sector)
|
|
|
|
|
if member_codes is None:
|
|
|
|
|
return self.store.get_latest_run(strategy, target_trade_date, query=query)
|
|
|
|
|
return self.store.get_latest_run(
|
|
|
|
|
strategy,
|
|
|
|
|
identity.target_trade_date,
|
|
|
|
|
query=replace(effective, sector=None),
|
|
|
|
|
sector_stock_codes=member_codes,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
def _sector_member_codes(
|
|
|
|
|
self,
|
|
|
|
|
target_trade_date: date,
|
|
|
|
|
sector_code: str,
|
|
|
|
|
) -> tuple[str, ...] | None:
|
|
|
|
|
"""Resolve one sector's members, or None when the port is absent."""
|
|
|
|
|
|
|
|
|
|
if self.sector_reader is None:
|
|
|
|
|
return None
|
|
|
|
|
return self.sector_reader.sector_member_codes(target_trade_date, sector_code)
|
|
|
|
|
|
|
|
|
|
def list_sector_counts(
|
|
|
|
|
self,
|
|
|
|
|
strategy: StrategyName,
|
|
|
|
|
target_trade_date: date | None = None,
|
|
|
|
|
*,
|
|
|
|
|
sector_type: str = "concept",
|
|
|
|
|
) -> SelectionSectorAggregates | None:
|
|
|
|
|
"""Aggregate the current run's selected stocks by point-in-time sector."""
|
|
|
|
|
|
|
|
|
|
run = self.store.get_latest_run(strategy, target_trade_date)
|
|
|
|
|
if run is None:
|
|
|
|
|
return None
|
|
|
|
|
selected_codes = [
|
|
|
|
|
item.ts_code
|
|
|
|
|
for item in run.items
|
|
|
|
|
if item.status == "selected" and item.signal_count > 0
|
|
|
|
|
]
|
|
|
|
|
membership = self._sector_membership(selected_codes, run.target_trade_date, sector_type)
|
|
|
|
|
return SelectionSectorAggregates(
|
|
|
|
|
run=run,
|
|
|
|
|
snapshot_trade_date=membership.snapshot_trade_date,
|
|
|
|
|
sector_type=sector_type,
|
|
|
|
|
sectors=membership.sector_counts,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
def _sector_membership(
|
|
|
|
|
self,
|
|
|
|
|
stock_codes: Sequence[str],
|
|
|
|
|
target_trade_date: date,
|
|
|
|
|
sector_type: str,
|
|
|
|
|
) -> SelectionSectorMembership:
|
|
|
|
|
"""Read sector counts for a stock set, tolerating a missing port."""
|
|
|
|
|
|
|
|
|
|
if self.sector_reader is None or not stock_codes:
|
|
|
|
|
return SelectionSectorMembership(snapshot_trade_date=None, sector_counts=())
|
|
|
|
|
return self.sector_reader.sector_counts(
|
|
|
|
|
stock_codes,
|
|
|
|
|
target_trade_date,
|
|
|
|
|
sector_type=sector_type,
|
|
|
|
|
)
|
2026-08-09 09:34:46 +08:00
|
|
|
|
|
|
|
|
|
2026-08-31 16:14:16 +08:00
|
|
|
def _to_item(
|
|
|
|
|
ts_code: str,
|
|
|
|
|
name: str,
|
|
|
|
|
evaluation: SelectionEvaluation,
|
|
|
|
|
*,
|
|
|
|
|
pattern_score: PatternScore | None = None,
|
|
|
|
|
) -> SelectionRunItem:
|
2026-08-09 09:34:46 +08:00
|
|
|
"""Translate a single-stock domain result into a stored item."""
|
|
|
|
|
|
|
|
|
|
return SelectionRunItem(
|
|
|
|
|
ts_code=ts_code,
|
|
|
|
|
name=name or (evaluation.signals[0].name if evaluation.signals else ""),
|
|
|
|
|
status=evaluation.status,
|
|
|
|
|
signal_count=len(evaluation.signals),
|
|
|
|
|
reason=evaluation.reason,
|
2026-08-31 16:14:16 +08:00
|
|
|
pattern_score=pattern_score or PatternScore(),
|
2026-08-09 09:34:46 +08:00
|
|
|
signals=evaluation.signals,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _run_status(evaluated_count: int, failed_count: int) -> SelectionRunStatus:
|
|
|
|
|
"""Map per-stock outcomes into a visible batch status."""
|
|
|
|
|
|
|
|
|
|
if failed_count == 0:
|
|
|
|
|
return "success"
|
|
|
|
|
if evaluated_count == 0 or failed_count >= evaluated_count:
|
|
|
|
|
return "failed"
|
|
|
|
|
return "partial_success"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _safe_item_error(error: Exception) -> str:
|
|
|
|
|
"""Keep per-stock failure context readable without persisting tracebacks."""
|
|
|
|
|
|
|
|
|
|
return " ".join(str(error).split())[:500] or error.__class__.__name__
|
|
|
|
|
|
|
|
|
|
|
2026-09-05 10:22:08 +08:00
|
|
|
def _turnover_present(history: StockHistory | None, target_trade_date: date) -> bool:
|
|
|
|
|
"""Return whether target-day Tushare turnover is available for diagnostics."""
|
|
|
|
|
|
|
|
|
|
if history is None:
|
|
|
|
|
return False
|
|
|
|
|
basic = history.daily_basic.get(target_trade_date)
|
|
|
|
|
return basic is not None and basic.turnover_rate is not None
|
|
|
|
|
|
|
|
|
|
|
2026-08-12 09:45:16 +08:00
|
|
|
def _chunks(
|
|
|
|
|
values: Sequence[SelectionStock],
|
|
|
|
|
size: int,
|
|
|
|
|
) -> tuple[tuple[SelectionStock, ...], ...]:
|
|
|
|
|
"""Split a stable stock sequence into bounded immutable chunks."""
|
|
|
|
|
|
|
|
|
|
return tuple(tuple(values[index : index + size]) for index in range(0, len(values), size))
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _chunk_count(value_count: int, size: int) -> int:
|
|
|
|
|
"""Return the number of chunks without materializing empty chunks."""
|
|
|
|
|
|
|
|
|
|
return (value_count + size - 1) // size
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _unique_stocks(stocks: Sequence[SelectionStock]) -> tuple[SelectionStock, ...]:
|
|
|
|
|
"""Keep the first source row for each stock so it is evaluated once."""
|
|
|
|
|
|
|
|
|
|
seen: set[str] = set()
|
|
|
|
|
unique: list[SelectionStock] = []
|
|
|
|
|
for stock in stocks:
|
|
|
|
|
ts_code = stock.ts_code
|
|
|
|
|
if ts_code in seen:
|
|
|
|
|
continue
|
|
|
|
|
seen.add(ts_code)
|
|
|
|
|
unique.append(stock)
|
|
|
|
|
return tuple(unique)
|
|
|
|
|
|
|
|
|
|
|
2026-08-09 09:34:46 +08:00
|
|
|
__all__ = [
|
|
|
|
|
"PreparedSelectionRun",
|
|
|
|
|
"RunZhixingB1",
|
|
|
|
|
"SelectionRerunRequired",
|
|
|
|
|
"SelectionRunInProgress",
|
2026-09-05 19:48:30 +08:00
|
|
|
"SelectionSectorAggregates",
|
2026-08-09 09:34:46 +08:00
|
|
|
]
|