Files
PolyWeather/main.py
T

578 lines
28 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import sys
import time
import os
import json
import re
from datetime import datetime, timedelta
from loguru import logger
from src.utils.config_loader import load_config
from src.utils.logger import setup_logger
from src.data_collection.polymarket_api import PolymarketClient
from src.data_collection.weather_sources import WeatherDataCollector
from src.data_collection.onchain_tracker import OnchainTracker
from src.models.statistical_model import TemperaturePredictor
from src.analysis.volume_analyzer import VolumeAnalyzer
from src.analysis.orderbook_analyzer import OrderbookAnalyzer
from src.analysis.technical_indicators import TechnicalIndicators
from src.analysis.whale_tracker import WhaleTracker
from src.strategy.decision_engine import DecisionEngine
from src.strategy.risk_manager import RiskManager
from src.trading.paper_trader import PaperTrader
from src.utils.notifier import TelegramNotifier
def main():
# 1. 初始化配置与日志
config_data = load_config()
setup_logger(config_data.get("app", {}).get("log_level", "INFO"))
logger.info("🌟 PolyWeather 监控引擎启动中...")
# 2. 初始化核心组件
polymarket = PolymarketClient(config_data["polymarket"])
weather = WeatherDataCollector(config_data["weather"])
onchain = OnchainTracker(config_data["polymarket"], polymarket)
notifier = TelegramNotifier(config_data["telegram"])
# 3. 初始化分析与交易组件
predictor = TemperaturePredictor()
decision_engine = DecisionEngine(config_data.get("config", {}))
whale_tracker = WhaleTracker(config_data.get("config", {}), onchain)
paper_trader = PaperTrader()
# 发送启动通知
notifier._send_message(
"🚀 <b>Polymarket 天气监控系统启动成功</b>\n正在扫描 12 个核心城市的最高温市场..."
)
# 信号记忆(持久化到文件)
pushed_signals = {}
SIGNALS_FILE = "data/pushed_signals.json"
if os.path.exists(SIGNALS_FILE):
try:
with open(SIGNALS_FILE, "r", encoding="utf-8") as f:
pushed_signals = json.load(f)
logger.info(f"已加载历史推送记录,共 {len(pushed_signals)} 条")
except:
pushed_signals = {}
# 确保data目录存在
if not os.path.exists("data"):
os.makedirs("data")
location_cache = {}
# 价格历史追踪(用于计算趋势)
PRICE_HISTORY_FILE = "data/price_history.json"
price_history = {}
if os.path.exists(PRICE_HISTORY_FILE):
try:
with open(PRICE_HISTORY_FILE, "r", encoding="utf-8") as f:
price_history = json.load(f)
except:
price_history = {}
try:
while True:
logger.info("--- 开启新一轮全量动态监控 (自动搜寻所有天气市场) ---")
cached_signals = {}
all_markets_cache = {}
# 1. 直接从 Polymarket 获取所有天气合约
all_weather_markets = polymarket.get_weather_markets()
# 1.5 尝试通过slug获取可能遗漏的市场(如部分结算的市场)
special_slugs = []
for slug in special_slugs:
event = polymarket.get_event_by_slug(slug)
if event:
title = event.get("title", "")
logger.info(f"通过slug找到特殊事件: {title}")
# 提取城市名
city = weather.extract_city_from_question(title)
if not city:
city = "Unknown"
# 将该事件的所有市场添加到列表
for m in event.get("markets", []):
# 检查是否已存在
c_id = m.get("conditionId")
if not any(
existing.get("condition_id") == c_id
for existing in all_weather_markets
):
all_weather_markets.append(
{
"condition_id": c_id,
"question": m.get("groupItemTitle")
or m.get("question"),
"active_token_id": m.get("activeTokenId"),
"tokens": m.get("clobTokenIds"),
"prices": m.get("outcomePrices"),
"event_title": title,
"slug": slug,
"city": city, # 提前标记城市
}
)
logger.debug(f"添加特殊市场: {m.get('groupItemTitle')}")
if not all_weather_markets:
logger.warning("当前 Polymarket 似乎没有任何活跃的天气市场,等待中...")
time.sleep(300)
continue
# 2. 批量同步盘口价格 (优化:为每个档位获取其对应的真实 Token 价格)
token_price_map = {}
price_requests = []
for m in all_weather_markets:
ts = m.get("tokens", [])
if isinstance(ts, str):
try: ts = json.loads(ts)
except: ts = []
active_tid = m.get("active_token_id")
# 如果是多选一市场(比如 Dallas 76-77°F
if len(ts) > 2 and active_tid:
# 获取该档位的买入价 (Ask)
price_requests.append({"token_id": active_tid, "side": "ask"})
# 获取该档位的买入“否”价所需的 Bid 价
price_requests.append({"token_id": active_tid, "side": "bid"})
# 如果是传统的 Yes/No 二选一市场
elif len(ts) == 2:
price_requests.append({"token_id": ts[0], "side": "ask"}) # Buy Yes
price_requests.append({"token_id": ts[1], "side": "ask"}) # Buy No
if price_requests:
logger.info(f"正在同步 {len(price_requests)} 个档位的真实盘口价格...")
token_price_map = polymarket.get_multiple_prices(price_requests)
logger.info(f"价格同步完成,成功获取 {len(token_price_map)} 个实时报价")
# 3. 按城市分组(按condition_id去重)
markets_by_city = {}
seen_condition_ids = set()
for i, m in enumerate(all_weather_markets):
c_id = m.get("condition_id")
if c_id in seen_condition_ids:
continue # 跳过重复
seen_condition_ids.add(c_id)
# 注入实时批量价格
ts = m.get("tokens", [])
if isinstance(ts, str):
try: ts = json.loads(ts)
except: ts = []
active_tid = m.get("active_token_id")
# 多选一市场逻辑
if len(ts) > 2 and active_tid:
m["buy_yes_live"] = token_price_map.get(active_tid) # 该档位的 Ask
# 买入“否”的价格 = 1 - 该档位的 Bid (别人愿意买 Yes 的最高价)
bid_price = token_price_map.get(active_tid) # 这里逻辑稍后在下文 fallback 处理中增强
# 我们这里只是注入,真正复杂的 fallback 逻辑在循环内部
# 二选一市场逻辑
elif len(ts) == 2:
m["buy_yes_live"] = token_price_map.get(ts[0])
m["buy_no_live"] = token_price_map.get(ts[1])
# 优先使用发现阶段已经识别出的城市名
city = m.get("city")
# 如果发现阶段没识别出,再尝试从问题文本提取
if not city or city == "Unknown":
full_context = f"{m.get('event_title', '')} {m.get('question', '')}"
city = weather.extract_city_from_question(full_context)
if i < 5:
logger.debug(
f"分析合约 {i}: City='{city}' | Title='{m.get('event_title')}"
)
if not city:
continue
if city not in markets_by_city:
markets_by_city[city] = []
markets_by_city[city].append(m)
logger.info(
f"动态发现 {len(markets_by_city)} 个受监控城市,共 {len(all_weather_markets)} 个合约"
)
# 3. 逐个城市分析
for city, city_markets in markets_by_city.items():
try:
# 获取/缓存坐标
if city not in location_cache:
coords = weather.get_coordinates(city)
if not coords:
continue
location_cache[city] = coords
logger.info(
f"📍 城市定位成功: {city} -> ({coords['lat']}, {coords['lon']})"
)
loc = location_cache[city]
# A. 获取实时天气共识
weather_data = weather.fetch_all_sources(
city, lat=loc["lat"], lon=loc["lon"]
)
consensus = weather.check_consensus(weather_data)
if not consensus.get("consensus"):
continue
logger.info(
f"☁️ {city} 当前气温: {consensus['average_temp']}°C | 监控合约: {len(city_markets)}"
)
# --- 本城市汇总预警缓存 ---
city_alerts = []
city_local_time = None
# B. 遍历该城市所有合约
for market in city_markets:
market_id = market.get("condition_id")
question = market.get("question", "未知市场")
event_title = market.get("event_title", "")
# 识别该合约的目标日期
target_date = weather.extract_date_from_title(
event_title
) or weather.extract_date_from_title(question)
ref_temp = consensus["average_temp"]
if target_date:
daily_data = weather_data.get("open-meteo", {}).get(
"daily", {}
)
if daily_data:
dates = daily_data.get("time", [])
max_temps = daily_data.get("temperature_2m_max", [])
for idx, d_str in enumerate(dates):
if target_date == d_str:
ref_temp = max_temps[idx]
break
# --- 价格获取逻辑 (增强版) ---
buy_yes_price = market.get("buy_yes_live")
buy_no_price = market.get("buy_no_live")
ts = market.get("tokens", [])
if isinstance(ts, str): ts = json.loads(ts)
active_tid = market.get("active_token_id")
# 特殊逻辑:针对多选一型市场(Dallas/Chicago 等温度段)
if len(ts) > 2 and active_tid:
# 1. 尝试从批量映射中重新获取该档位的精确买入价(Ask)
buy_yes_price = polymarket.get_price(active_tid, side="ask") if buy_yes_price is None else buy_yes_price
# 2. 计算“否”的价格(在多选一里是 1 - 最高买入意愿价Bid)
bid_price = polymarket.get_price(active_tid, side="bid")
if bid_price:
buy_no_price = 1.0 - bid_price
if i < 3: # 仅对前几个档位打印调试
logger.debug(f"Categorical价格匹配 [{market.get('city')}-{market.get('question')}]: Yes={buy_yes_price}, No={buy_no_price} (from Bid={bid_price})")
# 通用回退逻辑 (二选一)
if buy_yes_price is None or buy_no_price is None:
if ts and len(ts) >= 2:
prices = polymarket.get_buy_prices(ts[0], ts[1])
if prices:
buy_yes_price = prices.get("buy_yes")
buy_no_price = prices.get("buy_no")
# 最后的回退:使用原始 Gamma 概率
if buy_yes_price is None or buy_no_price is None:
gamma_prices = market.get("prices", [])
if isinstance(gamma_prices, str): gamma_prices = json.loads(gamma_prices)
prob = float(gamma_prices[0]) if gamma_prices else 0.5
if buy_yes_price is None: buy_yes_price = prob
if buy_no_price is None: buy_no_price = 1.0 - prob
# C. 准备缓存
temp_unit = weather_data.get("open-meteo", {}).get(
"unit", "celsius"
)
temp_symbol = "°F" if temp_unit == "fahrenheit" else "°C"
city_local_time = (
weather_data.get("open-meteo", {})
.get("current", {})
.get("local_time")
)
current_price = buy_yes_price if buy_yes_price else 0.5
# 计算价格趋势
prev_data = price_history.get(market_id, {})
prev_price = prev_data.get("price", current_price)
if prev_price > 0:
price_change_pct = ((current_price - prev_price) / prev_price) * 100
else:
price_change_pct = 0
# 更新价格历史
price_history[market_id] = {
"price": current_price,
"timestamp": datetime.now().isoformat()
}
cache_entry = {
"city": city,
"full_title": event_title,
"option": question,
"prediction": f"{ref_temp}{temp_symbol}",
"price": int(current_price * 100),
"buy_yes": int(buy_yes_price * 100),
"buy_no": int(buy_no_price * 100),
"url": f"https://polymarket.com/event/{market.get('slug')}",
"local_time": city_local_time,
"target_date": target_date,
"score": 0,
"rationale": "ACTIVE",
"trend": round(price_change_pct, 1),
}
if buy_yes_price <= 0.01 or buy_yes_price >= 0.99:
cache_entry["rationale"] = "ENDED"
all_markets_cache[market_id] = cache_entry
continue
# D. 评分
signal = decision_engine.calculate_signal(
model_prediction=predictor.predict_ensemble([ref_temp]),
market_data={
"orderbook": {},
"price_history": [current_price],
"transactions": [],
},
weather_consensus={"average_temp": ref_temp},
whale_activity=None,
)
cache_entry["score"] = signal["final_score"]
cache_entry["rationale"] = signal.get("recommendation", "N/A")
all_markets_cache[market_id] = cache_entry
# --- 预警收集 (仅监控价格) ---
if (0.85 <= buy_yes_price <= 0.95) or (
0.85 <= buy_no_price <= 0.95
):
alert_key = f"alert_{market_id}_range_85_95"
if alert_key not in pushed_signals:
trigger_side = (
"Buy Yes" if buy_yes_price >= 0.85 else "Buy No"
)
trigger_price = (
int(buy_yes_price * 100)
if trigger_side == "Buy Yes"
else int(buy_no_price * 100)
)
# --- 智能动态仓位计算 ---
# 1. 获取 Open-Meteo 对目标日期的最高温预测
predicted_high = None
weather_supports = False
daily_data = weather_data.get("open-meteo", {}).get("daily", {})
if daily_data and target_date:
dates = daily_data.get("time", [])
max_temps = daily_data.get("temperature_2m_max", [])
for idx, d_str in enumerate(dates):
if target_date == d_str and idx < len(max_temps):
predicted_high = max_temps[idx]
break
# 2. 判断天气预测是否支持当前方向
if predicted_high is not None:
# 解析选项的温度范围 (例如 "40-41°F" 或 "32°F or below")
temp_match = re.search(r'(\d+)(?:-(\d+))?°[FC]', question)
if temp_match:
low_bound = int(temp_match.group(1))
high_bound = int(temp_match.group(2)) if temp_match.group(2) else low_bound
# 如果买 NO,天气预测应该在这个区间之外
if trigger_side == "Buy No":
weather_supports = (predicted_high < low_bound - 2) or (predicted_high > high_bound + 2)
else: # 买 YES
weather_supports = (low_bound - 2 <= predicted_high <= high_bound + 2)
# 3. 获取成交量信息
market_volume = market.get("volume", 0)
if isinstance(market_volume, str):
try:
market_volume = float(market_volume.replace("$", "").replace(",", ""))
except:
market_volume = 0
high_volume = market_volume >= 5000 # $5000+ 算高成交量
# 4. 动态仓位决策
# 条件: 价格锁定程度 + 天气支持 + 成交量
if trigger_price >= 90 and weather_supports and high_volume:
# 三重确认:重注
amount_usd = 10.0
confidence_tag = "🔥高置信"
elif trigger_price >= 90 and weather_supports:
# 双重确认:中等仓位
amount_usd = 7.0
confidence_tag = "⭐中置信"
elif trigger_price >= 92:
# 价格接近锁定,即使其他条件不满足也小额参与
amount_usd = 5.0
confidence_tag = "📌价格锁定"
else:
# 普通信号:最小仓位
amount_usd = 3.0
confidence_tag = "💡试探"
logger.info(
f"【仓位决策】{city} {question} | "
f"价格:{trigger_price}¢ | 预测:{predicted_high} | 天气支持:{weather_supports} | "
f"高量:{high_volume} | 仓位:${amount_usd} ({confidence_tag})"
)
# --- 模拟交易触发逻辑 ---
side = "YES" if trigger_side == "Buy Yes" else "NO"
success = paper_trader.open_position(
market_id=market_id,
city=city,
option=question,
price=trigger_price,
side=side,
amount_usd=amount_usd,
target_date=target_date,
predicted_temp=predicted_high,
)
# 构建预测温度显示文本
temp_unit = weather_data.get("open-meteo", {}).get("unit", "celsius")
temp_symbol = "°F" if temp_unit == "fahrenheit" else "°C"
forecast_text = f"预测:{predicted_high}{temp_symbol}" if predicted_high else "预测:N/A"
city_alerts.append(
{
"type": "price",
"market": f"{question} ({target_date or '今日'})",
"msg": f"{trigger_side} {trigger_price}¢ | {forecast_text}",
"bought": success,
"amount": amount_usd,
"confidence": confidence_tag,
}
)
pushed_signals[alert_key] = time.time()
# 3. 信号暂存
cached_signals[market_id] = cache_entry
# --- 循环结束后统一推送本城市汇总 ---
if city_alerts:
notifier.send_combined_alert(
city, city_alerts, local_time=city_local_time
)
except Exception as e:
logger.error(f"分析城市 {city} 时出错: {e}")
continue
# --- 每处理完一个城市,立即更新 JSON 文件 ---
try:
# 1. 更新活跃信号缓存 (合并旧数据避免扫描中途变空)
final_signals = {}
if os.path.exists("data/active_signals.json"):
try:
with open(
"data/active_signals.json", "r", encoding="utf-8"
) as f:
final_signals = json.load(f)
except:
pass
final_signals.update(all_markets_cache)
with open("data/active_signals.json", "w", encoding="utf-8") as f:
json.dump(final_signals, f, ensure_ascii=False, indent=2)
# 2. 更新全量市场缓存
try:
with open("data/all_markets.json", "r", encoding="utf-8") as f:
existing_markets = json.load(f)
except:
existing_markets = {}
existing_markets.update(all_markets_cache)
# 清理过期日期
today_str = datetime.now().strftime("%Y-%m-%d")
cleaned_markets = {}
for k, v in existing_markets.items():
t_date = v.get("target_date")
if not t_date or t_date >= today_str:
cleaned_markets[k] = v
with open("data/all_markets.json", "w", encoding="utf-8") as f:
json.dump(cleaned_markets, f, ensure_ascii=False, indent=2)
# 3. 保存推送记录
with open("data/pushed_signals.json", "w", encoding="utf-8") as f:
json.dump(pushed_signals, f, ensure_ascii=False)
# 3.5 保存价格历史(用于趋势计算)
with open(PRICE_HISTORY_FILE, "w", encoding="utf-8") as f:
json.dump(price_history, f, ensure_ascii=False)
# --- 4. 更新模拟仓位盈亏 ---
price_snapshot = {}
for mid, entry in all_markets_cache.items():
price_snapshot[mid] = {"price": entry["price"]}
paper_trader.update_pnl(price_snapshot)
# --- 5. 每日收益总结推送 (北京时间 23:55 - 00:05 之间发送) ---
now_bj = datetime.utcnow() + timedelta(hours=8)
if now_bj.hour == 23 and now_bj.minute >= 50:
summary_key = f"daily_pnl_{now_bj.strftime('%Y%m%d')}"
if summary_key not in pushed_signals:
# 构造总结消息
total_cost = 0
total_pnl = 0
data = paper_trader._load_data()
pos_list = data.get("positions", {})
if pos_list:
report = [
f"📊 <b>每日模拟仓结算总结 ({now_bj.strftime('%Y-%m-%d')})</b>\n"
+ "═" * 15
]
for p in pos_list.values():
if p["status"] == "OPEN":
total_cost += p["cost_usd"]
total_pnl += p.get("pnl_usd", 0)
report.append(
f"💳 可用余额: <b>${data.get('balance', 0):.2f}</b>"
)
report.append(
f"💰 今日累计投入: <b>${total_cost:.2f}</b>"
)
report.append(
f"📈 累计浮动盈亏: <b>{total_pnl:+.2f}$</b>"
)
notifier._send_message("\n".join(report))
pushed_signals[summary_key] = time.time()
except Exception as e:
logger.error(f"即时保存数据失败: {e}")
logger.info("本轮扫描结束。等待 5 分钟...")
time.sleep(300)
except KeyboardInterrupt:
logger.info("收到关机指令,正在退出...")
except Exception as e:
logger.exception(f"系统运行出错: {e}")
if __name__ == "__main__":
main()