7a985dddd5
## 1. PG-only 重构
- 删 app/core/db/connection.py + schema.py (MySQL 路径)
- 新 app/core/db/{orm,models,pg_bootstrap}.py — SQLAlchemy 2.x ORM 一键建表
- 16 张业务表全在 market_data schema,原生 TIMESTAMPTZ / JSONB / Float / TEXT
- requirements.txt 删 PyMySQL 路径,加 psycopg2
## 2. 同步任务扩展(3 个新 task)
- **tick_trade** (mairui hsrl/zbjy):当天逐笔交易,21:00 发布
- **moneyflow** (mairui hsstock/history/transaction):个股资金流,21:30 发布
- **longhubang** (akshare):龙虎榜聚合层 + 席位层(2 张新表)
- 长虎榜放宽 akshare 政策:仅"无替代源 + 烟测通过"场景允许
- data_eastmoney 私有 API 不需要(akshare 烟测通过)
- 3-timer 设计:
- 15:30 market-sync.service (8 base tasks via runall_once)
- 21:05 market-sync-tick.service (tick_trade)
- 21:35 market-sync-moneyflow.service (moneyflow)
- 22:00 market-sync-lhb.service (longhubang,新加)
- bin/systemd/ 新增 tick / moneyflow / lhb 各 1 对 service+timer
- bin/market_sync_*_run.sh wrapper 脚本(不做法定节假日过滤,fail-open)
## 3. stock_code 统一为带 SH/SZ/BJ 前缀
- 历史 bug:stocks.code 用 SH600519,但 kline/moneyflow/tick_trade/kline_5min
/stock_sector_map/industry 6 张表用纯 6 位 600519,跨表 JOIN 全部 0 行
- 新增 to_hermes() 工具:6位 / 9位(mairui `000001.SZ` 格式)→ 统一 SH000001
- 5 个 task 改写:用 to_hermes(code6) 写入 stock_code
- 一次性迁移 6 张表存量 154M 行(CASE WHEN 探测 + 去重 + 加前缀)
- ORM: stock_sector_map.stock_code / industry.code String(6)→String(10)
## 4. Bug 修复
- **share table stock_code 格式**:之前写 6 位不带前缀,与 stocks 不一致
→ 修 task_share_snapshot + 一次性 UPDATE 63,417 行加前缀
- **share_snapshot warning 状态错填 last_error**:
→ 加 mark_sync_warning() 走专用路径,不写 last_failure_at / last_error
- **schedule config lastRun 不同步**:
→ 加 update_job_status_for_dataset(),SyncTask.run() 完成后自动镜像
→ cli/runall 触发的 task 也能更新 schedule config
## 5. mairui UTF-8 编码修复
- 历史 bug:mairui.py:_fetch 用 latin-1 兜底解码,把所有 UTF-8 中文名
double-encoded 写入 stocks.name(如 `歌华有线` 变成 `æ\xad\x8cå\x8d\x8e...`)
- 加 _decode_response():UTF-8 → GBK → latin-1 兜底
- 一次性修复 stocks.name 5,213 行:
- 4,370 行 (encode('latin-1').decode('utf-8') 反向解码)
- 616 行 (含 fullwidth A,宽松 printable 检查)
- 820 行 (mid-character 截断,重新从 mairui 拉)
## 6. 测试
- tests/test_smoke.py: TASKS 10→11, SYNC_DEFINITIONS 10→11
- tests/test_schema_models.py: 16→18 张表,新增 longhubang_daily/seat
- pytest 11/11 passed
## 验证
- 6 张表 0 残留无前缀行
- stocks JOIN kline_stock / kline_5min / moneyflow / tick_trade / stock_sector_map:88-100% 命中
- 5,213 stocks.name 全部正确 UTF-8 中文
- pytest 11/11 passed
136 lines
5.5 KiB
Python
136 lines
5.5 KiB
Python
"""CLI 入口:手动触发 / 状态查询。"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import json
|
||
import sys
|
||
from pathlib import Path
|
||
|
||
_PROJECT_ROOT = Path(__file__).resolve().parent.parent
|
||
if str(_PROJECT_ROOT) not in sys.path:
|
||
sys.path.insert(0, str(_PROJECT_ROOT))
|
||
|
||
|
||
def main():
|
||
parser = argparse.ArgumentParser(
|
||
description="market_data_sync CLI",
|
||
)
|
||
sub = parser.add_subparsers(dest="cmd", required=True)
|
||
|
||
# sync <task_id>
|
||
p_sync = sub.add_parser("sync", help="手动执行一次同步任务")
|
||
p_sync.add_argument("task_id", help="dataset_id, 例如 kline_daily")
|
||
p_sync.add_argument("--codes", help="指定股票代码,逗号分隔", default=None)
|
||
p_sync.add_argument("--start", help="起始日期 YYYY-MM-DD", default=None)
|
||
p_sync.add_argument("--end", help="结束日期 YYYY-MM-DD", default=None)
|
||
p_sync.add_argument("--workers", type=int, default=10, help="并发数")
|
||
p_sync.add_argument("--force", action="store_true", help="跳过时间门控(tick_trade 用)")
|
||
|
||
sub.add_parser("list", help="列出所有同步任务")
|
||
sub.add_parser("status", help="查看整体状态(调度器 / 数据源 / 同步任务)")
|
||
sub.add_parser("datasources", help="查看数据源健康度")
|
||
sub.add_parser("schedule", help="查看计划任务")
|
||
sub.add_parser("reseed", help="重新 seed 默认计划任务到 config 表")
|
||
|
||
args = parser.parse_args()
|
||
|
||
if args.cmd == "sync":
|
||
from app.core.utils.logging import setup_logging
|
||
setup_logging()
|
||
from app.tasks import get_task
|
||
from app.core.datasource.registry import build_default_registry
|
||
from app.core.sync.registry import seed_sync_registry
|
||
build_default_registry()
|
||
seed_sync_registry() # 确保 dataset_registry 有这条任务定义
|
||
try:
|
||
task = get_task(args.task_id)
|
||
except KeyError as e:
|
||
print(f"错误: {e}")
|
||
sys.exit(1)
|
||
kwargs = {}
|
||
if args.codes:
|
||
kwargs["codes"] = [c.strip() for c in args.codes.split(",") if c.strip()]
|
||
if args.start:
|
||
kwargs["start"] = args.start
|
||
if args.end:
|
||
kwargs["end"] = args.end
|
||
if args.workers:
|
||
kwargs["max_workers"] = args.workers
|
||
if args.force:
|
||
kwargs["force"] = True
|
||
result = task.run(trigger_source="cli", **kwargs)
|
||
print(json.dumps(result, ensure_ascii=False, indent=2))
|
||
sys.exit(0 if result.get("status") in ("ok", "warning") else 2)
|
||
|
||
if args.cmd == "list":
|
||
from app.tasks import TASKS
|
||
print("可用同步任务:")
|
||
for tid, cls in TASKS.items():
|
||
doc = (cls.__doc__ or cls.__name__).split("\n")[0]
|
||
print(f" {tid:25s} {doc}")
|
||
return
|
||
|
||
if args.cmd == "status":
|
||
from app.core.utils.logging import setup_logging
|
||
setup_logging()
|
||
from app.core.scheduler import get_scheduler_status
|
||
from app.core.datasource.registry import build_default_registry, get_health_status
|
||
from app.core.sync import get_registry_status
|
||
build_default_registry()
|
||
print("\n=== 调度器 ===")
|
||
s = get_scheduler_status()
|
||
print(json.dumps(s, ensure_ascii=False, indent=2, default=str))
|
||
print("\n=== 数据源健康度 ===")
|
||
for r in get_health_status():
|
||
mark = "✓" if r.get("success") else "✗"
|
||
print(f" {mark} {r.get('name', r.get('key')):30s} {r.get('message', '')}")
|
||
print("\n=== 同步任务状态 ===")
|
||
for r in get_registry_status():
|
||
print(f" {r['dataset_id']:25s} status={r.get('status'):10s} "
|
||
f"lastRun={r.get('last_success_at') or r.get('last_failure_at') or '-':20s} "
|
||
f"msg={r.get('message', '')[:60]}")
|
||
return
|
||
|
||
if args.cmd == "datasources":
|
||
from app.core.utils.logging import setup_logging
|
||
setup_logging()
|
||
from app.core.datasource.registry import build_default_registry, get_health_status, run_health_check
|
||
build_default_registry()
|
||
print("正在跑连通性测试...")
|
||
run_health_check()
|
||
for r in get_health_status():
|
||
mark = "✓" if r.get("success") else "✗"
|
||
print(f" {mark} {r.get('name', r.get('key')):30s} provides={r.get('provides', [])} {r.get('message', '')}")
|
||
return
|
||
|
||
if args.cmd == "schedule":
|
||
from app.core.utils.logging import setup_logging
|
||
setup_logging()
|
||
from app.core.scheduler import get_scheduler_status, seed_schedule_configs
|
||
seed_schedule_configs()
|
||
s = get_scheduler_status()
|
||
print(f"调度器运行中: {s['running']}")
|
||
for j in s["jobs"]:
|
||
enabled = "✓" if j.get("enabled") else "✗"
|
||
trigger = j.get("time") or f"every {j.get('intervalSeconds')}s"
|
||
print(f" {enabled} {j.get('name', j['key']):30s} {trigger:15s} "
|
||
f"condition={j.get('condition', 'always'):12s} "
|
||
f"last={j.get('lastStatus', '-')}")
|
||
return
|
||
|
||
if args.cmd == "reseed":
|
||
from app.core.utils.logging import setup_logging
|
||
setup_logging()
|
||
from app.core.scheduler.scheduler import seed_schedule_configs
|
||
from app.core.sync.registry import seed_sync_registry
|
||
from app.core.datasource.registry import seed_datasource_configs
|
||
seed_datasource_configs()
|
||
seed_sync_registry()
|
||
seed_schedule_configs()
|
||
print("已重新 seed:datasources + sync_registry + schedules")
|
||
return
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|