From be1afd248e4cc9b03a1e21090602a945af6988b2 Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 01:40:10 +0800 Subject: feat(app): 接入 AKShare 行情采集、净值查询与前端净值曲线 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit [变更性质] - 本提交完成 3b 行情闭环:真实 AKShare 搜索与净值采集、PostgreSQL 持久化、鉴权 API 与 Fable 前端净值曲线展示。 [新增功能] - API 新增 /api/instruments/search、/api/instruments/{code}/nav 与 nav/refresh 接口,经 Bearer 鉴权调用 AKShare Python 采集器并幂等落库。 - 前端提供基金搜索、精确代码选择、历史净值刷新/重读与 SVG 折线图,token 变更即清空私有结果并失效在途请求。 [实现方案] - MarketData 以 requiredProperty 校验采集器 payload,净值观测按 (code, nav_date) 幂等 upsert 并保留来源与哈希。 - 前端以 requestId 序列守卫 SearchCompleted/SearchFailed/NavCompleted/NavFailed,TokenChanged 同时递增 searchSeq/navSeq 拒绝过期响应;边界解码兼容 F# option 的 {"case":"Some"} 线格式。 [影响范围] - 新增 FundLab.Web.Tests(边界解码、序列失效、图表纯函数共 8 项),扩展 API 测试至 20 项;vite dev 代理 /api 至本地 API。 - 真实 AKShare 端到端依赖本机 akshare 环境,由负责人另行验收;本提交不包含 docs/overnight-progress.md 的现有修改。 --- src/FundLab.Api/akshare_collector.py | 171 +++++++++++++++++++++++++++++++++++ 1 file changed, 171 insertions(+) create mode 100644 src/FundLab.Api/akshare_collector.py (limited to 'src/FundLab.Api/akshare_collector.py') diff --git a/src/FundLab.Api/akshare_collector.py b/src/FundLab.Api/akshare_collector.py new file mode 100644 index 0000000..f988a61 --- /dev/null +++ b/src/FundLab.Api/akshare_collector.py @@ -0,0 +1,171 @@ +#!/usr/bin/env python3 +"""Small raw-data adapter for the AKShare endpoints used by Fund Lab.""" + +import argparse +import datetime as dt +import json +import re +import sys +from decimal import Decimal, InvalidOperation + +import akshare as ak +import pandas as pd + + +SCHEMA_VERSION = "fund-lab.akshare.v1" + + +def collected_at(): + return dt.datetime.now(dt.timezone.utc).isoformat().replace("+00:00", "Z") + + +def source_revision(): + return f"akshare-{getattr(ak, '__version__', 'unknown')}/eastmoney" + + +def is_missing(value): + try: + return bool(pd.isna(value)) + except (TypeError, ValueError): + return False + + +def text(value): + if is_missing(value): + return None + value = str(value).strip() + return value or None + + +def fund_code(value): + value = text(value) + if value is None: + return None + if value.isdigit() and len(value) < 6: + value = value.zfill(6) + return value if re.fullmatch(r"\d{6}", value) else None + + +def decimal_text(value): + value = text(value) + if value is None: + return None + try: + return format(Decimal(value), "f") + except InvalidOperation: + return None + + +def date_text(value): + if is_missing(value): + return None + if isinstance(value, (dt.datetime, dt.date)): + return value.strftime("%Y-%m-%d") + parsed = pd.to_datetime(value, errors="coerce") + if pd.isna(parsed): + return None + return parsed.strftime("%Y-%m-%d") + + +def search(query): + query = query.strip() + if not query: + raise ValueError("search query cannot be empty") + + frame = ak.fund_name_em() + query_lower = query.casefold() + rows = [] + + for _, row in frame.iterrows(): + code = fund_code(row.get("基金代码")) + name = text(row.get("基金简称")) + pinyin = text(row.get("拼音缩写")) + full_pinyin = text(row.get("拼音全称")) + + if code is None or name is None: + continue + + searchable = [code, name, pinyin or "", full_pinyin or ""] + if not any(query_lower in value.casefold() for value in searchable): + continue + + rows.append( + { + "code": code, + "name": name, + "fund_type": text(row.get("基金类型")), + "rank": 0 if code == query else (1 if name == query else 2), + } + ) + + rows.sort(key=lambda item: (item.pop("rank"), item["code"])) + return { + "schema_version": SCHEMA_VERSION, + "operation": "search", + "source": "akshare", + "source_revision": source_revision(), + "collected_at": collected_at(), + "instruments": rows[:100], + } + + +def nav(code): + if not re.fullmatch(r"\d{6}", code): + raise ValueError("fund code must contain exactly six digits") + + frame = ak.fund_open_fund_info_em( + symbol=code, + indicator="单位净值走势", + period="成立来", + ) + observations = [] + + for _, row in frame.iterrows(): + nav_date = date_text(row.get("净值日期")) + nav_value = decimal_text(row.get("单位净值")) + if nav_date is None or nav_value is None: + continue + + observations.append( + { + "nav_date": nav_date, + "published_at": None, + "nav": nav_value, + "accumulated_nav": None, + "daily_return": decimal_text(row.get("日增长率")), + } + ) + + if not observations: + raise ValueError(f"AKShare returned no usable NAV observations for {code}") + + return { + "schema_version": SCHEMA_VERSION, + "operation": "nav", + "source": "akshare", + "source_revision": source_revision(), + "collected_at": collected_at(), + "instrument": {"code": code}, + "observations": observations, + } + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--operation", choices=("search", "nav"), required=True) + parser.add_argument("--query") + parser.add_argument("--code") + args = parser.parse_args() + + try: + payload = search(args.query) if args.operation == "search" else nav(args.code) + json.dump(payload, sys.stdout, ensure_ascii=False, separators=(",", ":")) + sys.stdout.write("\n") + return 0 + except Exception as error: + print(f"AKShare collector failed: {error}", file=sys.stderr) + return 2 + + +if __name__ == "__main__": + raise SystemExit(main()) -- cgit v1.2.3