"""东方财富 API 客户端(资金流向数据)""" import asyncio import time import httpx import json import os import re from datetime import datetime, time as dtime, timedelta, timezone from typing import Optional, List from services.cache import get_cache, set_cache from services.cache import get_cache, set_cache # ---- API Key 轮询(MX 备选源用)---- _api_keys: List[str] = [] _key_index: int = -1 def _load_api_keys() -> List[str]: """收集所有可用 MX API key(支持逗号分隔和 MX_APIKEY_1 编号后缀)""" keys: List[str] = [] seen: set = set() def add(k: str): k = k.strip() if k and k not in seen: seen.add(k) keys.append(k) raw = os.environ.get("MX_APIKEY", "") if raw: for k in raw.split(","): add(k) for i in range(1, 10): raw = os.environ.get(f"MX_APIKEY_{i}", "") if raw: for k in raw.split(","): add(k) return keys def _ensure_keys() -> bool: global _api_keys, _key_index if not _api_keys: _api_keys = _load_api_keys() if _key_index == -1: val = get_cache("_mx_key_index") _key_index = int(val) if val else 0 return bool(_api_keys) def _get_key() -> str: """获取当前使用的 API key""" if not _ensure_keys(): return "" return _api_keys[_key_index % len(_api_keys)] def _rotate_key(): """切换到下一个 key 并持久化当前索引""" global _key_index _key_index += 1 set_cache("_mx_key_index", str(_key_index), ttl_hours=24) def parse_amount(text) -> float: """解析金额文本(支持 1.432亿元, -9915万元, 0元)""" if isinstance(text, (int, float)): return float(text) if not text or not isinstance(text, str): return 0 text = text.strip() if not text: return 0 sign = -1 if text.startswith("-") else 1 clean = text.lstrip("+-").strip() m = re.match(r"^([\d.]+)\s*(亿|万|元)?$", clean) if not m: return 0 value = float(m.group(1)) if m.group(1) else 0 unit = m.group(2) or "元" if unit == "亿": return sign * value * 100000000 elif unit == "万": return sign * value * 10000 return sign * value # ---- MX API(备选源)---- async def _call_mx_api(api_key: str, name: str, days: int) -> Optional[dict]: """执行一次 MX API 调用,返回原始 JSON""" url = "https://mkapi2.dfcfs.com/finskillshub/api/claw/query" payload = {"toolQuery": f"{name}最近{days}天主力资金流向和成交额"} async with httpx.AsyncClient() as client: try: resp = await client.post( url, json=payload, headers={ "Content-Type": "application/json", "apikey": api_key, }, timeout=15, ) if resp.status_code != 200: return None return resp.json() except Exception as e: print(f"[eastmoney] MX API error: {e}") return None def _parse_mx_response(result: dict) -> Optional[List[dict]]: """解析 MX API 返回,提取资金流向数据列表""" if result.get("status") != 0: return None dto_list = result.get("data", {}).get("data", {}).get("searchDataResultDTO", {}).get("dataTableDTOList", []) if not dto_list: return None dto = dto_list[0] name_map = dto.get("nameMap", {}) raw_table = dto.get("rawTable", {}) table = dto.get("table", {}) dates = table.get("headName", []) if not dates: return None name_to_id = {} for kid, v in name_map.items(): if isinstance(v, str): name_to_id[v] = kid main_net_col = name_to_id.get("(区间)主力净流入资金") amount_col = name_to_id.get("区间成交额") data_list = [] for i, date_str in enumerate(dates): date_str = str(date_str).strip() m1 = re.match(r"^(\d{4})[/-](\d{1,2})[/-](\d{1,2})", date_str) if m1: date_str = f"{m1.group(1)}-{int(m1.group(2)):02d}-{int(m1.group(3)):02d}" else: m2 = re.match(r"^(\d{1,2})[/-](\d{1,2})", date_str) if m2: year = datetime.now().year date_str = f"{year}-{int(m2.group(1)):02d}-{int(m2.group(2)):02d}" else: continue main_net = 0 if main_net_col and raw_table.get(main_net_col) and i < len(raw_table[main_net_col]): main_net = parse_amount(raw_table[main_net_col][i]) amount = 0 if amount_col and raw_table.get(amount_col) and i < len(raw_table[amount_col]): amount = parse_amount(raw_table[amount_col][i]) data_list.append({"date": date_str, "mainNetInflow": main_net, "amount": amount}) return data_list async def fetch_mx_api(name: str, days: int) -> Optional[List[dict]]: """调用东方财富妙想MX API获取资金流向(缓存6小时,多key轮询)""" cache_key = f"mx_fund_flow:{name}:{days}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) if not _ensure_keys(): return None for attempt in range(len(_api_keys)): key = _get_key() result = await _call_mx_api(key, name, days) if result is None: _rotate_key() continue status = result.get("status", -1) if status == 113: print(f"[eastmoney] key {_key_index % len(_api_keys)} 已达每日上限,切换到下一个") _rotate_key() continue if status == 114: print(f"[eastmoney] key {_key_index % len(_api_keys)} 无效(114),跳过") _rotate_key() continue if status != 0: _rotate_key() continue data_list = _parse_mx_response(result) if data_list: set_cache(cache_key, json.dumps(data_list, ensure_ascii=False)) return data_list print(f"[eastmoney] 所有 {len(_api_keys)} 个 MX API key 均已耗尽") return None def get_eastmoney_market(code: str) -> str: """获取东方财富格式的市场标识""" if code.startswith("688"): return "6" if code.startswith("60"): return "1" return "0" # ---- 板块数据 ---- # 从东方财富 bkzj/list.js 逆向的字段映射 # f62=主力净流入, f184=主力净流入占比 # f66=超大单净流入, f69=超大单净流入占比 # f72=大单净流入, f75=大单净流入占比 # f78=中单净流入, f81=中单净流入占比 # f84=小单净流入, f87=小单净流入占比 # f70=成交额 SECTOR_FIELDS = "f12,f14,f2,f3,f62,f184,f66,f69,f72,f75,f78,f81,f84,f87,f70" # 东方财富板块类型映射 SECTOR_MEDIA_MAP = { "industry": "m:90+s:4", "concept": "m:90+t:3", } # 东方财富 UT 令牌管理 _em_ut: str = "8dec03ba335b81bf4ebdf7b29ec27d15" _em_ut_lock = asyncio.Lock() async def _refresh_em_ut() -> str: """ 从东方财富前端 JS 中提取最新的 ut 令牌。 按优先级尝试: 1. bkzj/list.js(板块页专用) 2. common/emdataview.js(通用数据组件) """ urls = [ "https://data.eastmoney.com/newstatic/js/bkzj/list.js", "https://data.eastmoney.com/newstatic/js/common/emdataview.js", ] headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", "Referer": "https://data.eastmoney.com/bkzj/hy.html", } async with httpx.AsyncClient() as client: for url in urls: try: resp = await client.get(url, headers=headers, timeout=10) if resp.status_code != 200: continue # 匹配 ut: 'xxxx' 或 ut:'xxxx' 或 ut: "xxxx" m = re.search(r"""ut['"]?\s*:\s*['"]([a-f0-9]{32})['"]""", resp.text) if m: token = m.group(1) print(f"[eastmoney] 已刷新 UT 令牌: {token[:8]}...") return token except Exception as e: print(f"[eastmoney] 获取 UT 失败({url}): {e}") return _em_ut # 保底返回当前值 async def get_em_ut(force_refresh: bool = False) -> str: """获取当前 UT,必要时刷新""" global _em_ut if force_refresh: async with _em_ut_lock: _em_ut = await _refresh_em_ut() return _em_ut # ---- 板块数据(市场时间感知缓存)---- _CST = timezone(timedelta(hours=8)) # 北京时间 _TRADING_MORNING = (dtime(9, 30), dtime(11, 30)) _TRADING_AFTERNOON = (dtime(13, 0), dtime(15, 0)) def _cst_now() -> datetime: return datetime.now(_CST) def _is_trading_time() -> bool: """判断当前是否为 A 股交易时段(周一至周五 9:30-11:30 / 13:00-15:00)""" now = _cst_now() if now.weekday() >= 5: return False t = now.time() return (_TRADING_MORNING[0] <= t <= _TRADING_MORNING[1] or _TRADING_AFTERNOON[0] <= t <= _TRADING_AFTERNOON[1]) def _sector_ttl_hours() -> int: """根据是否在交易时段返回缓存 TTL - 交易时段: 2 分钟(数据持续变化) - 非交易时段: 18 小时(覆盖到下一个交易日) """ return 0 if _is_trading_time() else 18 # curl_cffi 模拟 Chrome TLS 指纹 from curl_cffi.requests import AsyncSession _sector_session: Optional[AsyncSession] = None def _get_sector_session() -> AsyncSession: global _sector_session if _sector_session is None: _sector_session = AsyncSession( impersonate="chrome120", headers={ "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", "Referer": "https://data.eastmoney.com/bkzj/hy.html", "Accept": "*/*", "Accept-Language": "zh-CN,zh;q=0.9", }, timeout=10, ) return _sector_session # 内存缓存(加速交易时段频繁请求) _sector_cache: dict[str, tuple[list[dict], float]] = {} _SECTOR_MEM_TTL = 60 async def fetch_sector_list(sector_type: str) -> list[dict]: # 1. 内存缓存(仅交易时段有效,60s 避免重复请求) now = time.time() if sector_type in _sector_cache: data, ts = _sector_cache[sector_type] if now - ts < _SECTOR_MEM_TTL: return data cache_key = f"sector_list:{sector_type}" # 2. 非交易时段:走磁盘持久缓存,不请求 API if not _is_trading_time(): cached = get_cache(cache_key) if cached is not None: data = json.loads(cached) _sector_cache[sector_type] = (data, now) return data # 3. 请求 API data = await _fetch_push2(sector_type) if data: _sector_cache[sector_type] = (data, now) ttl = _sector_ttl_hours() # 交易时段 ttl=0 → 不写磁盘(内存缓存已够) if ttl > 0: set_cache(cache_key, json.dumps(data, ensure_ascii=False), ttl_hours=ttl) return data data = await _fetch_akshare(sector_type) if data: _sector_cache[sector_type] = (data, now) # akshare 数据源更新不确定,不写入磁盘缓存 return data async def _fetch_push2(sector_type: str) -> list[dict]: """东方财富 push2 API(curl_cffi 模拟浏览器 TLS 指纹)""" fs = SECTOR_MEDIA_MAP.get(sector_type) if not fs: return [] session = _get_sector_session() ut = await get_em_ut() url = ( f"https://push2.eastmoney.com/api/qt/clist/get" f"?fs={fs}&fields={SECTOR_FIELDS}" f"&fid=f62&po=1&pz=500&pn=1&np=1&fltt=2" f"&invt=2&ut={ut}" ) for attempt in range(2): try: resp = await session.get(url) if resp.status_code != 200: if attempt == 0: await asyncio.sleep(1) continue return [] result = resp.json() if result.get("rc") != 0: return [] diff = result.get("data", {}).get("diff", []) items = [] for item in diff: items.append({ "code": item.get("f12", ""), "name": item.get("f14", ""), "level": item.get("f2"), "changePercent": item.get("f3"), "changeAmount": None, "mainNetInflow": item.get("f62", 0) or 0, "mainNetInflowPercent": item.get("f184", 0), "superLargeInflow": item.get("f66", 0) or 0, "superLargeInflowPercent": item.get("f69", 0), "largeInflow": item.get("f72", 0) or 0, "largeInflowPercent": item.get("f75", 0), "mediumInflow": item.get("f78", 0) or 0, "mediumInflowPercent": item.get("f81", 0), "smallInflow": item.get("f84", 0) or 0, "smallInflowPercent": item.get("f87", 0), "turnover": item.get("f70", 0) or 0, }) return items except Exception as e: err = str(e) print(f"[eastmoney] push2 获取{sector_type}板块失败(attempt {attempt+1}): {err[:80]}") # UT 可能过期,尝试刷新 if "disconnect" in err.lower() or "refused" in err.lower() or attempt == 1: await get_em_ut(force_refresh=True) ut = _em_ut # 重建会话(TLS 指纹可能会被缓存) global _sector_session _sector_session = None session = _get_sector_session() if attempt == 0: await asyncio.sleep(1) return [] import akshare as ak async def _fetch_akshare(sector_type: str) -> list[dict]: """akshare 降级方案(同花顺数据源)""" loop = asyncio.get_event_loop() def _get(): code_map = {} try: if sector_type == "industry": code_df = ak.stock_board_industry_name_em() else: code_df = ak.stock_board_concept_name_em() if code_df is not None and not code_df.empty: for _, r in code_df.iterrows(): code_map[str(r.get("f14", ""))] = str(r.get("f12", "")) except Exception: try: if sector_type == "industry": code_df = ak.stock_board_industry_name_ths() else: code_df = ak.stock_board_concept_name_ths() if code_df is not None and not code_df.empty: for _, r in code_df.iterrows(): code_map[str(r.get("name", ""))] = str(r.get("code", "")) except Exception: pass if sector_type == "industry": df = ak.stock_fund_flow_industry() else: df = ak.stock_fund_flow_concept() return code_map, df try: code_map, df = await loop.run_in_executor(None, _get) if df is None or df.empty: return [] df = df.sort_values("净额", ascending=False) items = [] for _, row in df.iterrows(): name = str(row.get("行业", "")).strip() inflow = float(row.get("流入资金", 0) or 0) * 100000000 outflow = float(row.get("流出资金", 0) or 0) * 100000000 items.append({ "code": code_map.get(name, ""), "name": name, "level": float(row.get("行业指数") or 0), "changePercent": float(row.get("行业-涨跌幅") or 0), "changeAmount": None, "mainNetInflow": float(row.get("净额", 0) or 0) * 100000000, "mainNetInflowPercent": None, "superLargeInflow": None, "superLargeInflowPercent": None, "largeInflow": None, "largeInflowPercent": None, "mediumInflow": None, "mediumInflowPercent": None, "smallInflow": None, "smallInflowPercent": None, "turnover": inflow + outflow, }) return items except Exception as e: print(f"[eastmoney] akshare 获取{sector_type}板块失败: {e}") return [] # ---- 公司概况 ---- _F10_MARKET_MAP = {"6": "SH", "0": "SZ", "3": "SZ"} _DATACENTER_MARKET_MAP = {"6": "SH", "0": "SZ", "3": "SZ"} def _code_to_f10(code: str) -> str: """6位代码 → 东方财富 F10 格式 (SH600519 / SZ000001)""" prefix = code[0] mkt = _F10_MARKET_MAP.get(prefix, "SH") return f"{mkt}{code}" def _code_to_secucode(code: str) -> str: """6位代码 → 数据中心格式 (600519.SH / 000001.SZ)""" prefix = code[0] mkt = _DATACENTER_MARKET_MAP.get(prefix, "SH") return f"{code}.{mkt}" async def fetch_company_profile(code: str) -> Optional[dict]: """东方财富 F10 - 公司概况(缓存24小时)""" cache_key = f"company_profile:{code}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) em_code = _code_to_f10(code) url = ( f"https://emweb.securities.eastmoney.com/PC_HSF10" f"/CompanySurvey/CompanySurveyAjax?code={em_code}" ) headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": ( f"https://emweb.securities.eastmoney.com/PC_HSF10" f"/CompanySurvey/Index?type=web&code={em_code}" ), } async with httpx.AsyncClient() as client: try: resp = await client.get(url, headers=headers, timeout=15) if resp.status_code != 200: return None raw = resp.json() except Exception as e: print(f"[eastmoney] fetch_company_profile error: {e}") return None jbzl = raw.get("jbzl", {}) fxxg = raw.get("fxxg", {}) profile = { "code": code, "name": jbzl.get("agjc", ""), "fullName": jbzl.get("gsmc", ""), "englishName": jbzl.get("ywmc", ""), "boardType": jbzl.get("zqlb", ""), "industry": jbzl.get("sshy", ""), "exchange": jbzl.get("ssjys", ""), "industryDetail": jbzl.get("sszjhhy", ""), # 高管 "chairman": jbzl.get("dsz", ""), "generalManager": jbzl.get("zjl", ""), "secretary": jbzl.get("dm", ""), # 联系方式 "phone": jbzl.get("lxdh", ""), "email": jbzl.get("dzxx", ""), "fax": jbzl.get("cz", ""), "website": jbzl.get("gswz", ""), # 地址 "officeAddress": jbzl.get("bgdz", ""), "registeredAddress": jbzl.get("zcdz", ""), "region": jbzl.get("qy", ""), "postalCode": jbzl.get("yzbm", ""), # 注册信息 "registeredCapital": jbzl.get("zczb", ""), "unifiedCreditCode": jbzl.get("gsdj", ""), "employees": jbzl.get("gyrs", 0), "managementCount": jbzl.get("glryrs", 0), # 中介机构 "lawFirm": jbzl.get("lssws", ""), "auditor": jbzl.get("kjssws", ""), # 公司简介&业务 "businessScope": jbzl.get("jyfw", ""), "companyProfile": jbzl.get("gsjj", ""), # 发行信息 "establishmentDate": fxxg.get("clrq", ""), "listingDate": fxxg.get("ssrq", ""), "listingPrice": fxxg.get("mgfxj", ""), "issuePriceEarningsRatio": fxxg.get("fxsyl", ""), "issueAmount": fxxg.get("fxl", ""), "netProceeds": fxxg.get("mjzjje", ""), "issueMethod": fxxg.get("fxfs", ""), } set_cache(cache_key, json.dumps(profile, ensure_ascii=False), ttl_hours=24) return profile # ---- 财务数据 ---- # 财务指标关键字段映射(英文 → 中文) FINANCIAL_FIELD_LABELS: dict[str, str] = { "EPSJB": "基本每股收益", "EPSKCJB": "扣非每股收益", "BPS": "每股净资产", "MGWFPLR": "每股未分配利润", "MGJYXJJE": "每股经营现金流", "TOTALOPERATEREVE": "营业总收入", "MLR": "毛利润", "PARENTNETPROFIT": "归母净利润", "KCFJCXSYJLR": "扣非净利润", "TOTALOPERATEREVETZ": "营收同比增长", "PARENTNETPROFITTZ": "净利润同比增长", "KCFJCXSYJLRTZ": "扣非净利润同比增长", "ROEJQ": "加权净资产收益率", "ROEKCJQ": "扣非ROE", "ZZCJLL": "总资产净利率", "XSJLL": "销售净利率", "XSMLL": "销售毛利率", "LD": "流动比率", "SD": "速动比率", "ZCFZL": "资产负债率", "XJLLB": "现金流比率", "YSZKZZTS": "应收账款周转天数", "CHZZTS": "存货周转天数", "ZZCZZTS": "总资产周转天数", "TOAZZL": "总资产周转率", "CHZZL": "存货周转率", "ROIC": "投入资本回报率", "QYCS": "权益乘数", "CQBL": "产权比率", "TAXRATE": "实际税率", "TOTAL_SHARE": "总股本", "A_FREE_SHARE": "流通A股", "DJD_TOI_YOY": "单季营收同比", "DJD_DPNP_YOY": "单季净利润同比", "EPSJBTZ": "每股收益同比增长", "BPSTZ": "每股净资产同比增长", "ROEJQTZ": "ROE同比增长", # 利润表补充 "OPERATE_PROFIT_PK": "营业利润", "PER_TOI": "每股营业收入", "PER_EBIT": "每股EBIT", # 资产负债表补充 "TOTAL_ASSETS_PK": "总资产", "LIABILITY": "总负债", "TOTAL_EQUITY_PK": "所有者权益", "MGZBGJ": "每股资本公积", "CA_TA": "流动资产占比", "NCA_TA": "非流动资产占比", "CL_TL": "流动负债占比", # 现金流量表补充 "NETCASH_OPERATE_PK": "经营活动现金流净额", "NETCASH_INVEST_PK": "投资活动现金流净额", "NETCASH_FINANCE_PK": "筹资活动现金流净额", "JYXJLYYSR": "经营现金流占营收比", "XSJXLYYSR": "销售现金流占营收比", "FCFF_FORWARD": "企业自由现金流", # 营运能力补充 "OPERATE_CYCLE": "营业周期", "CURRENT_ASSET_TR": "流动资产周转率", "FIXED_ASSET_TR": "固定资产周转率", "ACCOUNTS_PAYABLE_TR": "应付账款周转率", "PREPAID_ACCOUNTS_RATIO": "预收账款占营收比", "PAYABLE_TDAYS": "应付账款周转天数", } # 小数比率字段(后端返回小数,前端需显示为百分比) _DECIMAL_RATIO_FIELDS = { "JYXJLYYSR", "XSJXLYYSR", "PREPAID_ACCOUNTS_RATIO", "NCO_OP", "NCO_FIXED", "SALE_ER_PK", } def _parse_financial_record(record: dict) -> dict: """将原始财务记录转为整洁的输出""" result = {} for field, label in FINANCIAL_FIELD_LABELS.items(): val = record.get(field) if val is not None: if isinstance(val, (int, float)): # 小数比率字段转为百分比 if field in _DECIMAL_RATIO_FIELDS and abs(val) < 1: result[label] = round(val * 100, 2) result[label + "(%)"] = True elif abs(val) > 1_0000_0000: result[label] = round(val / 1_0000_0000, 2) result[label + "(亿)"] = True elif abs(val) > 1_0000: result[label] = round(val / 1_0000, 2) result[label + "(万)"] = True else: result[label] = round(val, 2) else: result[label] = val return result async def fetch_financial_data_main(code: str, years: int = 5) -> Optional[list[dict]]: """东方财富数据中心 - 主要财务指标(缓存6小时)""" cache_key = f"financial_main:{code}:{years}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) secucode = _code_to_secucode(code) url = "https://datacenter.eastmoney.com/securities/api/data/v1/get" params = { "reportName": "RPT_F10_FINANCE_MAINFINADATA", "columns": "ALL", "filter": f'(SECUCODE="{secucode}")', "pageNumber": 1, "pageSize": max(years * 4, 4), # 每年最多4份报告 "sortTypes": -1, "sortColumns": "REPORT_DATE", } headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": "https://datacenter.eastmoney.com/", } async with httpx.AsyncClient() as client: try: resp = await client.get(url, params=params, headers=headers, timeout=15) if resp.status_code != 200: return None body = resp.json() if not body.get("result") or not body["result"].get("data"): return None records = body["result"]["data"] except Exception as e: print(f"[eastmoney] fetch_financial_data error: {e}") return None items = [] for r in records: items.append({ "reportDate": r.get("REPORT_DATE", ""), "reportType": r.get("REPORT_TYPE", ""), "reportDateName": r.get("REPORT_DATE_NAME", ""), "noticeDate": r.get("NOTICE_DATE", ""), "indicators": _parse_financial_record(r), }) set_cache(cache_key, json.dumps(items, ensure_ascii=False), ttl_hours=6) return items # ---- 经营分析(主营构成/收入构成)---- _SEGMENT_TYPE_MAP = { "industry": "1", # 按行业 "product": "2", # 按产品 "region": "3", # 按地区 } async def fetch_business_segments(code: str, by_type: str = "product", years: int = 3) -> Optional[list[dict]]: """东方财富经营分析 - 主营构成(按产品/行业/地区),缓存6小时 Args: code: 6位股票代码 by_type: "product"|"industry"|"region" years: 返回最近N年的数据 """ seg_type = _SEGMENT_TYPE_MAP.get(by_type, "2") cache_key = f"business_segments:{code}:{by_type}:{years}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) secucode = _code_to_secucode(code) url = "https://datacenter.eastmoney.com/securities/api/data/v1/get" params = { "reportName": "RPT_F10_FN_SEGMENTSV", "columns": ( "SECUCODE,SECURITY_CODE,SECURITY_NAME_ABBR,MAINOP_TYPE,ITEM_NAME," "ITEM_CODE,ITEM_PARENT_CODE,MAIN_BUSINESS_INCOME,MAIN_BUSINESS_RPOFIT," "GROSS_RPOFIT_RATIO,MBI_RATIO,MBR_RATIO,REPORT_DATE,REPORT_NAME," "REPORT_DATE_CODE,MAIN_BUSINESS_COST,MBC_RATIO" ), "filter": f'(SECUCODE="{secucode}")(MAINOP_TYPE="{seg_type}")(ITEM_PARENT_CODE="999999")', "pageNumber": 1, "pageSize": 200, "sortTypes": "1,1", "sortColumns": "MAINOP_TYPE,RANK3", } headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": "https://datacenter.eastmoney.com/", } async with httpx.AsyncClient() as client: try: resp = await client.get(url, params=params, headers=headers, timeout=15) if resp.status_code != 200: return None body = resp.json() if not body.get("result") or not body["result"].get("data"): return None records = body["result"]["data"] except Exception as e: print(f"[eastmoney] fetch_business_segments error: {e}") return None # 按报告期分组 from collections import OrderedDict period_map: dict[str, list[dict]] = OrderedDict() for r in records: report_name = r.get("REPORT_NAME", "") report_date = r.get("REPORT_DATE", "") item_name = r.get("ITEM_NAME", "") income = r.get("MAIN_BUSINESS_INCOME") cost = r.get("MAIN_BUSINESS_COST") profit = r.get("MAIN_BUSINESS_RPOFIT") gross_ratio = r.get("GROSS_RPOFIT_RATIO") mbi_ratio = r.get("MBI_RATIO") mbr_ratio = r.get("MBR_RATIO") mbc_ratio = r.get("MBC_RATIO") entry = { "itemName": item_name, "income": income, "cost": cost, "profit": profit, "grossRatio": round(gross_ratio * 100, 2) if gross_ratio is not None else None, "incomeRatio": round(mbi_ratio * 100, 2) if mbi_ratio is not None else None, "profitRatio": round(mbr_ratio * 100, 2) if mbr_ratio is not None else None, "costRatio": round(mbc_ratio * 100, 2) if mbc_ratio is not None else None, } if report_name not in period_map: period_map[report_name] = [] period_map[report_name].append(entry) # 补充 profitRatio:若 API 未返回,用该期总利润推算各项目利润占比 for pname, items in period_map.items(): total_profit = sum(it["profit"] for it in items if it["profit"] is not None) if total_profit != 0: for it in items: if it["profitRatio"] is None and it["profit"] is not None: it["profitRatio"] = round(it["profit"] / total_profit * 100, 2) # 取最近N年 all_periods = list(period_map.keys()) all_periods.sort(reverse=True) selected = all_periods[:max(years * 4, 4)] result = [] for p in selected: result.append({ "reportName": p, "items": period_map[p], }) set_cache(cache_key, json.dumps(result, ensure_ascii=False), ttl_hours=6) return result async def fetch_kline_history(code: str, days: int = 90) -> Optional[list[dict]]: """从 push2his kline/get 获取日K线(curl_cffi 模拟浏览器 TLS 指纹) 可能受 IP 限流影响,调用方应有降级。 返回: [{date, open, close, high, low, volume(手), turnover, changePercent, amplitude, turnoverRate}] """ market = get_eastmoney_market(code) secid = f"{market}.{code}" url = ( f"https://push2his.eastmoney.com/api/qt/stock/kline/get" f"?secid={secid}" f"&ut=b2884a393a59ad64002292a3e90d46a5" f"&fields1=f1,f2,f3,f4,f5,f6" f"&fields2=f51,f52,f53,f54,f55,f56,f57,f58,f59,f60,f61" f"&klt=101&fqt=1&end=20500101&lmt={days}" ) from curl_cffi.requests import AsyncSession async with AsyncSession(impersonate="chrome120") as session: for attempt in range(2): try: resp = await session.get(url, headers={ "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": "https://quote.eastmoney.com/", }, timeout=10) if resp.status_code != 200: if attempt == 0: continue return None result = resp.json() if result.get("rc") != 0: if attempt == 0: continue return None klines = result.get("data", {}).get("klines", []) if not klines: if attempt == 0: continue return None parsed = [] for line in klines: fields = line.split(",") if len(fields) < 11: continue def _f(i): return float(fields[i]) if fields[i] else 0 parsed.append({ "date": fields[0], "open": _f(1), "close": _f(2), "high": _f(3), "low": _f(4), "volume": int(_f(5)), # 手 "turnover": _f(6), "amplitude": _f(7), "changePercent": _f(8), "change": _f(9), "turnoverRate": _f(10), }) parsed.reverse() return parsed except Exception as e: if attempt == 0: continue return None return None async def fetch_fund_flow_daykline(code: str, days: int) -> Optional[list[dict]]: """从 push2his daykline/get 获取资金流向日线数据(结构化返回,替代 MX API) 返回字段:date, mainNetInflow, superLargeInflow/Outflow, largeInflow/Outflow, mediumInflow/Outflow, smallInflow/Outflow, 各占比%, closePrice, changePercent, mainForceNet, retailNet """ market = get_eastmoney_market(code) secid = f"{market}.{code}" url = ( f"https://push2his.eastmoney.com/api/qt/stock/fflow/daykline/get" f"?lmt={days}&fields1=f1,f2,f3,f7" f"&fields2=f51,f52,f53,f54,f55,f56,f57,f58,f59,f60,f61,f62,f63" f"&ut=b2884a393a59ad64002292a3e90d46a5&secid={secid}" ) headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": "https://quote.eastmoney.com/", } for attempt in range(2): async with httpx.AsyncClient() as client: try: resp = await client.get(url, headers=headers, timeout=10) if resp.status_code != 200: if attempt == 0: continue return None result = resp.json() klines = result.get("data", {}).get("klines", []) if not klines: if attempt == 0: continue return None parsed = [] for line in klines: fields = line.split(",") if len(fields) < 13: continue def _f(i): return float(fields[i]) if fields[i] else 0 # 字段索引: 0=f51(date), 1=f52(主力), 2=f53(小单), 3=f54(中单), # 4=f55(大单), 5=f56(超大单), 6=f57(主力%), 7=f58(小单%), # 8=f59(中单%), 9=f60(大单%), 10=f61(超大单%), 11=f62(收盘价), 12=f63(涨跌幅) date = fields[0] main_net = _f(1) small_net = _f(2) medium_net = _f(3) large_net = _f(4) super_large_net = _f(5) parsed.append({ "date": date, "mainNetInflow": main_net, "superLargeInflow": max(0, super_large_net), "superLargeOutflow": abs(min(0, super_large_net)), "largeInflow": max(0, large_net), "largeOutflow": abs(min(0, large_net)), "mediumInflow": max(0, medium_net), "mediumOutflow": abs(min(0, medium_net)), "smallInflow": max(0, small_net), "smallOutflow": abs(min(0, small_net)), "mainNetInflowPercent": _f(6), "superLargeInflowPercent": _f(10), "superLargeOutflowPercent": 0, "largeInflowPercent": _f(9), "largeOutflowPercent": 0, "mediumInflowPercent": _f(8), "mediumOutflowPercent": 0, "smallInflowPercent": _f(7), "smallOutflowPercent": 0, "turnover": 0, # daykline 不含成交额,由调用方从 Tencent K 线合并 "mainForceNet": main_net, "retailNet": small_net + medium_net, "changePercent": _f(12), "closePrice": _f(11), }) return parsed except Exception as e: if attempt == 0: continue print(f"[eastmoney] fetch_fund_flow_daykline error: {e}") return None return None