| import asyncio |
| import logging |
| import random |
| from collections import defaultdict |
| from datetime import datetime, timedelta, timezone |
|
|
| import numpy as np |
|
|
| from app.core.contracts import ( |
| CATEGORIES, |
| CONTRACTS, |
| EXCHANGES, |
| ContractInfo, |
| get_contracts_by_category, |
| get_contracts_by_exchange, |
| ) |
| from app.models.schemas import KlineData, MarketData |
|
|
| logger = logging.getLogger(__name__) |
|
|
| _INTERVAL_MINUTES = {"1m": 1, "5m": 5, "15m": 15, "30m": 30, "1h": 60, "4h": 240, "1d": 1440, "1w": 10080} |
| _BAR_COUNT = {"1m": 500, "5m": 300, "15m": 200, "30m": 200, "1h": 200, "4h": 150, "1d": 120, "1w": 60} |
| _VOLATILITY = {"1m": 0.0004, "5m": 0.0009, "15m": 0.0015, "30m": 0.002, "1h": 0.003, "4h": 0.005, "1d": 0.01, "1w": 0.02} |
|
|
|
|
| def _now(): |
| return datetime.now(timezone.utc) |
|
|
|
|
| class MarketDataService: |
| INTERVALS = list(_INTERVAL_MINUTES.keys()) |
|
|
| def __init__(self): |
| self._quotes: dict[str, MarketData] = {} |
| self._klines: dict[str, dict[str, list[KlineData]]] = defaultdict(lambda: defaultdict(list)) |
| self._subscribers: list[asyncio.Queue] = [] |
| self._running = False |
| self._task: asyncio.Task | None = None |
| self._mode: str = "simulated" |
| self._real_data_cache: dict[str, dict] = {} |
| self._sim_prices: dict[str, float] = {} |
| self._init_simulated() |
|
|
| @property |
| def mode(self): |
| return self._mode |
|
|
| def set_mode(self, mode: str): |
| if mode not in ("simulated", "realtime"): |
| raise ValueError("mode must be 'simulated' or 'realtime'") |
| self._mode = mode |
|
|
| def get_exchanges(self): |
| return EXCHANGES |
|
|
| def get_categories(self): |
| return CATEGORIES |
|
|
| def get_contracts(self, category=None, exchange=None): |
| cs = get_contracts_by_category(category) if category else get_contracts_by_exchange(exchange) if exchange else CONTRACTS |
| return [{"symbol": c.symbol, "exchange": c.exchange, "category": c.category, "name_cn": c.name_cn} for c in cs] |
|
|
| def get_current_price(self, symbol): |
| q = self._quotes.get(symbol) |
| return q.close if q else self._sim_prices.get(symbol) |
|
|
| def get_quote(self, symbol): |
| return self._quotes.get(symbol) |
|
|
| def get_all_market_data(self): |
| return list(self._quotes.values()) |
|
|
| def get_kline_history(self, symbol, interval="1m", limit=200): |
| return self._klines.get(symbol, {}).get(interval, [])[-limit:] |
|
|
| def subscribe(self): |
| q: asyncio.Queue = asyncio.Queue(maxsize=200) |
| self._subscribers.append(q) |
| return q |
|
|
| def unsubscribe(self, q): |
| if q in self._subscribers: |
| self._subscribers.remove(q) |
|
|
| async def _notify(self, data): |
| dead = [] |
| for q in self._subscribers: |
| try: |
| q.put_nowait(data) |
| except asyncio.QueueFull: |
| dead.append(q) |
| for d in dead: |
| self._subscribers.remove(d) |
|
|
| async def start(self): |
| if self._running: |
| return |
| self._running = True |
| self._task = asyncio.create_task(self._run()) |
| logger.info("MarketDataService started (%s, %d contracts)", self._mode, len(CONTRACTS)) |
|
|
| async def stop(self): |
| self._running = False |
| if self._task: |
| self._task.cancel() |
| try: |
| await self._task |
| except asyncio.CancelledError: |
| pass |
|
|
| async def _run(self): |
| tick = 0 |
| while self._running: |
| try: |
| if self._mode == "realtime" and tick % 5 == 0: |
| await self._fetch_realtime() |
| else: |
| await self._tick_sim() |
| except Exception as e: |
| logger.error("Tick error: %s", e) |
| tick += 1 |
| await asyncio.sleep(1) |
|
|
| |
|
|
| async def _fetch_realtime(self): |
| loop = asyncio.get_event_loop() |
| for c in CONTRACTS: |
| try: |
| df = await loop.run_in_executor(None, self._fetch_one, c.symbol) |
| if df is not None and len(df) > 0: |
| r = df.iloc[0] |
| md = MarketData( |
| symbol=c.symbol, name=c.name_cn, exchange=c.exchange, category=c.category, |
| open=float(r.get("open", 0)), high=float(r.get("high", 0)), |
| low=float(r.get("low", 0)), |
| close=float(r.get("trade", 0) or r.get("close", 0)), |
| volume=int(r.get("volume", 0)), |
| bid=float(r.get("bidprice1", 0) or 0), ask=float(r.get("askprice1", 0) or 0), |
| open_interest=int(r.get("position", 0) or 0), |
| pre_close=float(r.get("preclose", 0) or 0), |
| pre_settlement=float(r.get("prevsettlement", 0) or r.get("presettlement", 0) or 0), |
| settlement=float(r.get("settlement", 0) or 0), |
| change_pct=float(r.get("changepercent", 0) or 0), |
| timestamp=_now(), |
| ) |
| self._quotes[c.symbol] = md |
| await self._notify(md) |
| self._real_data_cache[c.symbol] = {"contracts": df.to_dict("records"), "updated_at": _now().isoformat()} |
| except Exception as e: |
| logger.debug("Fetch %s: %s", c.symbol, e) |
|
|
| @staticmethod |
| def _fetch_one(symbol): |
| import akshare as ak |
| return ak.futures_zh_realtime(symbol=symbol) |
|
|
| def get_contract_details(self, symbol): |
| return self._real_data_cache.get(symbol) |
|
|
| |
|
|
| _BASE = { |
| "沪深300指数期货": 3800, "上证50指数期货": 2500, "中证500指数期货": 5600, |
| "中证1000股指期货": 5200, "5年期国债期货": 102.5, "10年期国债期货": 104.2, |
| "2年期国债期货": 101.3, |
| "黄金": 530, "白银": 7200, "铂": 980, "钯": 1050, |
| "沪铜": 72000, "沪铝": 19500, "沪锌": 22000, "沪铅": 16000, |
| "沪镍": 128000, "沪锡": 250000, "不锈钢": 13500, "国际铜": 65000, |
| "氧化铝": 3800, "铸造铝合金期货": 20000, "工业硅": 10500, "碳酸锂": 76000, |
| "多晶硅": 42000, |
| "螺纹钢": 3500, "热轧卷板": 3600, "线材": 3700, "铁矿石": 850, |
| "焦煤": 1350, "焦炭": 2100, "锰硅": 6500, "硅铁": 7200, "动力煤": 850, |
| "原油": 550, "燃油": 3200, "沥青": 3600, "橡胶": 14000, "20号胶": 11000, |
| "丁二烯橡胶": 12000, "纸浆": 5200, "胶版印刷纸期货": 5500, "低硫燃料油": 3500, |
| "PTA": 5800, "郑醇": 2500, "乙二醇": 4500, "尿素": 1700, "纯碱": 1500, |
| "玻璃": 1400, "塑料": 8000, "PP": 7500, "PVC": 5600, "苯乙烯": 8200, |
| "短纤": 7000, "液化石油气": 4200, "纯苯": 6800, "烧碱": 3200, |
| "二甲苯": 8500, "瓶级聚酯切片": 6200, "丙烯": 7500, |
| "豆一": 4200, "豆二": 3600, "豆粕": 3000, "豆油": 7800, "棕榈": 7200, |
| "菜油": 8500, "菜籽": 5200, "菜粕": 2300, "花生": 8500, |
| "棉花": 14000, "棉纱": 22000, "白糖": 6500, "鲜苹果": 8200, "红枣": 10000, |
| "玉米": 2200, "玉米淀粉": 2500, "鸡蛋": 3500, "生猪": 14000, |
| "粳米": 3200, "粳稻": 2700, "强麦": 2600, "早籼稻": 2500, "晚籼稻": 2600, |
| "纤维板": 1200, "胶合板": 160, "原木": 1800, |
| "集运指数(欧线)期货": 1200, |
| } |
|
|
| def _init_simulated(self): |
| for c in CONTRACTS: |
| base = self._BASE.get(c.symbol, 5000.0) |
| self._generate_contract(c, base) |
|
|
| def _generate_contract(self, contract: ContractInfo, base: float): |
| sym = contract.symbol |
| now = _now() |
|
|
| for interval in self.INTERVALS: |
| count = _BAR_COUNT[interval] |
| vol = _VOLATILITY[interval] |
| mins = _INTERVAL_MINUTES[interval] |
| price = base * (0.9 + np.random.random() * 0.1) |
| bars: list[KlineData] = [] |
|
|
| for i in range(count): |
| ts = now - timedelta(minutes=mins * (count - i)) |
| change = np.random.normal(0, vol) |
| op = round(price, 2) |
| cl = round(price * (1 + change), 2) |
| hi = round(max(op, cl) * (1 + abs(np.random.normal(0, vol * 0.4))), 2) |
| lo = round(min(op, cl) * (1 - abs(np.random.normal(0, vol * 0.4))), 2) |
| vl = int(np.random.exponential(300 * mins) + 50) |
| bars.append(KlineData(symbol=sym, interval=interval, open=op, high=hi, low=lo, close=cl, volume=vl, timestamp=ts)) |
| price = cl |
|
|
| self._klines[sym][interval] = bars |
|
|
| final_price = self._klines[sym]["1m"][-1].close |
| self._sim_prices[sym] = final_price |
|
|
| day_bars = self._klines[sym]["1d"] |
| recent = day_bars[-1] if day_bars else None |
| md = MarketData( |
| symbol=sym, name=contract.name_cn, exchange=contract.exchange, category=contract.category, |
| open=recent.open if recent else base, |
| high=recent.high if recent else base, |
| low=recent.low if recent else base, |
| close=final_price, |
| volume=recent.volume if recent else 0, |
| bid=round(final_price - final_price * 0.00005, 2), |
| ask=round(final_price + final_price * 0.00005, 2), |
| pre_close=round(base, 2), |
| change_pct=round((final_price - base) / base, 6), |
| timestamp=_now(), |
| ) |
| self._quotes[sym] = md |
|
|
| async def _tick_sim(self): |
| vol = 0.0003 |
| for c in CONTRACTS: |
| sym = c.symbol |
| price = self._sim_prices.get(sym, 5000.0) |
| change = np.random.normal(0, vol) |
| np_ = round(price * (1 + change), 2) |
| self._sim_prices[sym] = np_ |
|
|
| hi = round(max(price, np_) * (1 + abs(random.gauss(0, vol * 0.2))), 2) |
| lo = round(min(price, np_) * (1 - abs(random.gauss(0, vol * 0.2))), 2) |
| vl = int(np.random.exponential(200) + 30) |
|
|
| kline = KlineData(symbol=sym, interval="1m", open=price, high=hi, low=lo, close=np_, volume=vl, timestamp=_now()) |
| ks = self._klines[sym]["1m"] |
| ks.append(kline) |
| if len(ks) > 2000: |
| self._klines[sym]["1m"] = ks[-1500:] |
|
|
| prev = self._quotes.get(sym) |
| md = MarketData( |
| symbol=sym, name=c.name_cn, exchange=c.exchange, category=c.category, |
| open=prev.open if prev else price, |
| high=round(max(hi, prev.high if prev else hi), 2), |
| low=round(min(lo, prev.low if prev else lo), 2), |
| close=np_, |
| volume=(prev.volume if prev else 0) + vl, |
| bid=round(np_ - np_ * 0.00005, 2), ask=round(np_ + np_ * 0.00005, 2), |
| pre_close=round(prev.pre_close if prev else price, 2), |
| change_pct=round((np_ - (prev.pre_close if prev else price)) / (prev.pre_close if prev else price), 6), |
| timestamp=_now(), |
| ) |
| self._quotes[sym] = md |
| await self._notify(md) |
|
|
|
|
| market_data_service = MarketDataService() |
|
|