feat(selection): implement gold brick resonance strategy with evaluation and logging enhancements

This commit is contained in:
yuxuanhui
2026-09-05 10:22:08 +08:00
parent 9163590070
commit 79476252b1
20 changed files with 1060 additions and 97 deletions
@@ -1,16 +1,17 @@
"""Application orchestration for persisted whole-universe B1 runs."""
"""Application orchestration for persisted whole-universe strategy runs."""
from __future__ import annotations
import logging
import time
from collections.abc import Callable, Sequence
from collections import Counter
from collections.abc import Callable, Mapping, Sequence
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass
from datetime import date
from typing import Literal, Protocol, cast
from typing import Protocol, cast
from ..domain.models import SelectionEvaluation, StockHistory
from ..domain.models import SelectionEvaluation, SelectionStrategyName, StockHistory
from ..domain.pattern_scoring import (
PatternCase,
PatternCaseLibraryLoader,
@@ -33,9 +34,16 @@ from ..domain.runs import (
)
from .evaluate import EvaluateZhixingB1
logger = logging.getLogger(__name__)
StrategyName = Literal["zhixing_b1"]
_FAILURE_STATUSES = {"insufficient_history", "missing_target_bar", "data_error"}
# 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",
}
class SelectionEvaluator(Protocol):
@@ -53,13 +61,18 @@ class PreparedSelectionRun:
class RunZhixingB1:
"""Prepare, execute, and query persisted Zhixing B1 result batches."""
"""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.
"""
def __init__(
self,
reader: SelectionUniverseReader,
store: SelectionRunStore,
evaluator: SelectionEvaluator | None = None,
evaluators: Mapping[StrategyName, SelectionEvaluator] | None = None,
pattern_case_loader: PatternCaseLibraryLoader | None = None,
pattern_scorer: PatternScorer | None = None,
*,
@@ -76,6 +89,11 @@ class RunZhixingB1:
self.reader = reader
self.store = store
self.evaluator = evaluator or EvaluateZhixingB1(reader)
self.evaluators: dict[StrategyName, SelectionEvaluator] = {
"zhixing_b1": self.evaluator,
}
if evaluators is not None:
self.evaluators.update(evaluators)
self.pattern_case_loader = pattern_case_loader
self.pattern_scorer = pattern_scorer
self.pattern_scoring_enabled = pattern_scoring_enabled
@@ -91,12 +109,43 @@ class RunZhixingB1:
) -> PreparedSelectionRun:
"""Validate source eligibility before claiming the rerunnable key."""
source = self.reader.load_execution_source(strategy, target_trade_date)
run = self.store.prepare_run(
logger.info(
"selection_run_prepare_started strategy=%s target_trade_date=%s rerun=%s",
strategy,
target_trade_date,
source,
rerun=rerun,
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,
)
return PreparedSelectionRun(run=run, source=source)
@@ -109,29 +158,76 @@ class RunZhixingB1:
"""
stocks = _unique_stocks(prepared.source.stocks)
strategy = prepared.run.strategy
target_trade_date = prepared.source.target_trade_date
evaluator = self.evaluators[strategy]
evaluated_count = 0
selected_stock_count = 0
signal_count = 0
failed_count = 0
missing_turnover_count = 0
insufficient_history_count = 0
history_rows = 0
batch_count = _chunk_count(len(stocks), self.batch_size)
current_batch = 0
final_status: SelectionRunStatus = "failed"
read_seconds = 0.0
evaluate_seconds = 0.0
persist_seconds = 0.0
scoring_seconds = 0.0
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,
)
try:
pattern_cases, pattern_library_error = self._prepare_pattern_cases(prepared.run.id)
if strategy == "zhixing_b1":
pattern_cases, pattern_library_error = self._prepare_pattern_cases(prepared.run.id)
else:
pattern_cases, pattern_library_error = None, None
with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
for batch_stocks in _chunks(stocks, self.batch_size):
for batch_index, batch_stocks in enumerate(
_chunks(stocks, self.batch_size),
start=1,
):
current_batch = batch_index
read_started = time.perf_counter()
histories = self._load_histories(
batch_stocks,
prepared.source.target_trade_date,
target_trade_date,
evaluator,
)
read_seconds += time.perf_counter() - read_started
history_rows += sum(
batch_read_seconds = time.perf_counter() - read_started
read_seconds += batch_read_seconds
batch_history_rows = sum(
len(history.bars) for history in histories if history is not None
)
history_rows += batch_history_rows
batch_missing_turnover = sum(
not _turnover_present(history, target_trade_date)
for history in histories
)
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,
)
evaluate_started = time.perf_counter()
evaluations = tuple(
@@ -139,10 +235,46 @@ class RunZhixingB1:
self._evaluate_stock,
batch_stocks,
histories,
[prepared.source.target_trade_date] * len(batch_stocks),
[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),
)
)
evaluate_seconds += time.perf_counter() - evaluate_started
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,
)
scoring_started = time.perf_counter()
items = tuple(
@@ -151,6 +283,7 @@ class RunZhixingB1:
stock.name,
evaluation,
pattern_score=self._score_stock(
strategy,
prepared.run.id,
stock,
history,
@@ -173,23 +306,79 @@ class RunZhixingB1:
signal_count += sum(item.signal_count for item in items)
failed_count += sum(item.status in _FAILURE_STATUSES for item in items)
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,
)
persist_started = time.perf_counter()
self._record_items(prepared.run.id, items)
persist_seconds += time.perf_counter() - persist_started
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,
)
status = _run_status(evaluated_count, failed_count)
final_status = _run_status(evaluated_count, failed_count)
self.store.finish_run(
prepared.run.id,
status,
final_status,
evaluated_count=evaluated_count,
selected_stock_count=selected_stock_count,
signal_count=signal_count,
failed_count=failed_count,
)
except Exception as exc: # noqa: BLE001 - worker boundary must persist failure state
logger.error(
"selection_run_failed run_id=%s error_type=%s reason=%s",
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,
)
except Exception as exc: # noqa: BLE001 - worker boundary must persist failure state
final_status = "failed"
logger.error(
"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(),
prepared.run.id,
current_batch,
exc.__class__.__name__,
_safe_item_error(exc),
)
@@ -202,23 +391,41 @@ class RunZhixingB1:
signal_count=signal_count,
failed_count=max(failed_count, 1),
error_type="batch_error",
error_message=str(exc),
error_message=_safe_item_error(exc),
)
except Exception: # noqa: BLE001 - preserve the original worker failure
logger.error(
"selection_run_failure_persist_failed run_id=%s",
"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(),
prepared.run.id,
current_batch,
)
finally:
logger.info(
"selection_run_summary run_id=%s stock_count=%d history_rows=%d "
"batch_count=%d worker_count=%d read_seconds=%.3f "
"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 "
"evaluate_seconds=%.3f scoring_seconds=%.3f persist_seconds=%.3f",
strategy,
target_trade_date.isoformat(),
prepared.run.id,
prepared.source.market_sync_batch_id,
len(stocks),
history_rows,
batch_count,
current_batch,
self.max_workers,
evaluated_count,
selected_stock_count,
signal_count,
failed_count,
insufficient_history_count,
missing_turnover_count,
final_status,
read_seconds,
evaluate_seconds,
scoring_seconds,
@@ -251,6 +458,7 @@ class RunZhixingB1:
def _score_stock(
self,
strategy: StrategyName,
run_id: str,
stock: SelectionStock,
history: StockHistory | None,
@@ -260,7 +468,11 @@ class RunZhixingB1:
) -> PatternScore:
"""Score one selected stock once and isolate enrichment failures."""
if not self.pattern_scoring_enabled or evaluation.status != "selected":
if (
strategy != "zhixing_b1"
or not self.pattern_scoring_enabled
or evaluation.status != "selected"
):
return PatternScore()
if library_error is not None:
return PatternScore.failed(library_error)
@@ -283,6 +495,7 @@ class RunZhixingB1:
self,
stocks: Sequence[SelectionStock],
target_trade_date: date,
evaluator: SelectionEvaluator,
) -> tuple[StockHistory | None, ...]:
"""Load one chunk when the reader supports it, with old-path fallback."""
@@ -300,7 +513,8 @@ class RunZhixingB1:
for stock in typed_stocks
)
if isinstance(self.evaluator, EvaluateZhixingB1):
execute_history = getattr(evaluator, "execute_history", None)
if callable(execute_history):
return tuple(
self.reader.load_history(stock.ts_code, target_trade_date) for stock in typed_stocks
)
@@ -311,23 +525,35 @@ class RunZhixingB1:
stock: SelectionStock,
history: StockHistory | None,
target_trade_date: date,
evaluator: SelectionEvaluator,
strategy: StrategyName,
run_id: str,
batch_index: int,
) -> 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(
self.evaluator,
evaluator,
"execute_history",
None,
)
if history is not None and execute_history is not None:
return execute_history(history, target_trade_date)
return self.evaluator.execute(ts_code, target_trade_date)
return evaluator.execute(ts_code, target_trade_date)
except Exception as exc: # noqa: BLE001 - isolate one stock from the batch
logger.warning(
"selection_item_failed ts_code=%s error_type=%s reason=%s",
"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,
ts_code,
len(history.bars) if history is not None else 0,
_turnover_present(history, target_trade_date),
exc.__class__.__name__,
_safe_item_error(exc),
)
@@ -407,6 +633,15 @@ def _safe_item_error(error: Exception) -> str:
return " ".join(str(error).split())[:500] or error.__class__.__name__
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
def _chunks(
values: Sequence[SelectionStock],
size: int,