summaryrefslogtreecommitdiff
path: root/src/FundLab.Api/akshare_collector.py
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-21 01:40:10 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-21 01:40:10 +0800
commitbe1afd248e4cc9b03a1e21090602a945af6988b2 (patch)
treed31bd6e6c1736ab4ad784b52ebc5010ade074439 /src/FundLab.Api/akshare_collector.py
parent2ada70d6467aec11b45328112d598454ca50f2d6 (diff)
downloadfund-lab-be1afd248e4cc9b03a1e21090602a945af6988b2.tar.gz
feat(app): 接入 AKShare 行情采集、净值查询与前端净值曲线
[变更性质] - 本提交完成 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 的现有修改。
Diffstat (limited to 'src/FundLab.Api/akshare_collector.py')
-rw-r--r--src/FundLab.Api/akshare_collector.py171
1 files changed, 171 insertions, 0 deletions
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())