"""东方财富题材数据服务:题材列表、题材详情、题材相关股票 逆向自 emrnweb.eastmoney.com/investment 的 H5 接口: - 题材列表: POST https://emcfgdata.eastmoney.com/api/themeInvest/getThemeList - 题材详情: GET https://emcfgdata.securities.eastmoney.com/api/themeInvest/getDetail/{themeCode} - 相关股票: POST https://emcfgdata.eastmoney.com/api/themeInvest/getStockList 两个 POST 接口需要「移动端包装结构」: {args:{...业务参数}, appKey, client, clientVersion, clientType, randomCode, timestamp} """ import asyncio import json import random import string import time from datetime import datetime, time as dtime, timedelta, timezone from typing import Optional import httpx from services.cache import get_cache, set_cache # ---- 东方财富移动端配置中心域名 ---- _PZ_URL = "https://emcfgdata.eastmoney.com" _PZ_CDN_URL = "https://emcfgdata.securities.eastmoney.com" # appKey 与页面场景对应:题材列表索引页 / 题材详情页 _APP_KEY_INDEX = "rn-themeIndex" _APP_KEY_DETAIL = "rn-themeDetail" _HEADERS = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "Origin": "https://emrnweb.eastmoney.com", "Referer": "https://emrnweb.eastmoney.com/", "Accept": "application/json", "Accept-Language": "zh-CN,zh;q=0.9", } # ---- 交易时段感知缓存 ---- _CST = timezone(timedelta(hours=8)) # 北京时间 _TRADING_MORNING = (dtime(9, 30), dtime(11, 30)) _TRADING_AFTERNOON = (dtime(13, 0), dtime(15, 0)) def _is_trading_time() -> bool: """判断当前是否为 A 股交易时段(周一至周五 9:30-11:30 / 13:00-15:00)""" now = datetime.now(_CST) 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 _dynamic_ttl() -> int: """盘中返回 2 分钟缓存 TTL,非交易时段 18 小时(覆盖到下一交易日)""" return 0 if _is_trading_time() else 18 # ---- 请求封装 ---- def _build_payload(args: Optional[dict] = None, app_key: str = _APP_KEY_INDEX) -> dict: """构建东方财富移动端请求包装结构""" return { "args": args or {}, "appKey": app_key, "client": "iOS", "clientVersion": "8.3", "clientType": "cfw", "randomCode": "".join(random.choices(string.ascii_uppercase + string.ascii_lowercase + string.digits, k=16)), "timestamp": int(time.time() * 1000), } async def _post(path: str, args: dict, app_key: str = _APP_KEY_INDEX) -> Optional[dict]: """POST 到配置中心接口,返回 data 层 JSON""" payload = _build_payload(args, app_key) try: async with httpx.AsyncClient(timeout=15) as client: resp = await client.post(_PZ_URL + path, json=payload, headers=_HEADERS) if resp.status_code != 200: return None body = resp.json() except Exception as e: print(f"[themes] POST {path} 失败: {e}") return None if body.get("code") != 0: print(f"[themes] POST {path} 返回错误: {body.get('message')}") return None return body.get("data") async def _get_cdn(path: str, app_key: str = _APP_KEY_DETAIL) -> Optional[dict]: """GET 到配置中心 CDN 接口(题材详情),data 包装结构放 query 参数""" payload = _build_payload({}, app_key) params = {"data": json.dumps(payload, ensure_ascii=False)} try: async with httpx.AsyncClient(timeout=15) as client: resp = await client.get(_PZ_CDN_URL + path, params=params, headers=_HEADERS) if resp.status_code != 200: return None body = resp.json() except Exception as e: print(f"[themes] GET {path} 失败: {e}") return None if body.get("code") != 0: print(f"[themes] GET {path} 返回错误: {body.get('message')}") return None return body.get("data") # ---- 题材列表 ---- # sortField 映射(题材列表页):1=涨幅(bf3) 3=强度(strengthValue) 4=热度排名(hotRank) 5=成交额(fex5) _LIST_PAGE_SIZE = 500 async def fetch_theme_list(sort_field: int = 1, asc: bool = False) -> list[dict]: """获取全部题材列表(内部循环分页拉全,约 623 个,最多 2 页) Args: sort_field: 排序字段 1/3/4/5 asc: True=升序, False=降序 """ cache_key = f"theme_list:{sort_field}:{asc}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) sort = 1 if asc else -1 # hotRank 数值越小越热,"热度降序(最热在前)" 需反转为接口升序 if sort_field == 4: sort = -sort items: list[dict] = [] page = 1 total = None for _ in range(5): # 安全上限 data = await _post( "/api/themeInvest/getThemeList", {"pageSize": _LIST_PAGE_SIZE, "pageNum": page, "sort": sort, "sortField": sort_field}, ) if not data: break if total is None: total = data.get("total", 0) page_items = data.get("list", []) if not page_items: break items.extend(page_items) if len(items) >= total: break page += 1 if items: ttl = _dynamic_ttl() if ttl > 0: set_cache(cache_key, json.dumps(items, ensure_ascii=False), ttl_hours=ttl) return items # ---- 题材详情 ---- async def fetch_theme_detail(theme_code: str) -> Optional[dict]: """获取题材详情(简介 + 热点事件 + 相关新闻),缓存 1 小时""" cache_key = f"theme_detail:{theme_code}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) data = await _get_cdn(f"/api/themeInvest/getDetail/{theme_code}", app_key=_APP_KEY_DETAIL) if not data: return None set_cache(cache_key, json.dumps(data, ensure_ascii=False), ttl_hours=1) return data # ---- 题材相关股票 ---- _STOCK_PAGE_SIZE = 100 async def fetch_theme_stocks(theme_code: str) -> dict: """获取题材下全部相关股票(分页拉全),返回 {stockList, statistic, total}""" cache_key = f"theme_stocks:{theme_code}" cached = get_cache(cache_key) if cached is not None: return json.loads(cached) stock_list: list[dict] = [] statistic = {} total = 0 page = 1 for _ in range(20): # 安全上限 data = await _post( "/api/themeInvest/getStockList", {"themeCode": theme_code, "pageSize": _STOCK_PAGE_SIZE, "pageNum": page, "sort": -1, "sortField": "f3"}, app_key=_APP_KEY_DETAIL, ) if not data: break if not statistic and data.get("statistic"): statistic = data["statistic"] page_items = data.get("stockList", []) if not page_items: break stock_list.extend(page_items) total = data.get("total", 0) if len(stock_list) >= total: break page += 1 result = {"stockList": stock_list, "statistic": statistic, "total": total} if stock_list: ttl = _dynamic_ttl() if ttl > 0: set_cache(cache_key, json.dumps(result, ensure_ascii=False), ttl_hours=ttl) return result