Files
stock/main.py
T

178 lines
7.9 KiB
Python

import os
import sys
import sqlite3
import logging
import configparser
from datetime import datetime, timedelta
import pandas as pd
import numpy as np
import requests
from quant_strategy import QuantStateMachine
from util import get_history_k, get_5min_k
# =====================================================================
# 初始化基础环境与日志
# =====================================================================
current_dir = os.path.dirname(os.path.abspath(__file__))
log_filename = os.path.join(current_dir, "quant_daily.log")
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s',
handlers=[logging.FileHandler(log_filename, encoding='utf-8')])
config = configparser.ConfigParser()
config.read(os.path.join(current_dir, "config.ini"), encoding='utf-8')
BOT_TOKEN = config.get("telegram", "bot_token")
CHAT_ID = config.get("telegram", "chat_id")
DB_FILE = os.path.join(current_dir, config.get("database", "db_name", fallback="watchlist.db"))
PROXY_ENABLED = config.getint("proxy", "enabled", fallback=0)
PROXY_URL = config.get("proxy", "url", fallback="")
# =====================================================================
# 补全缺失的函数 1: 真实/模拟数据接口 (请在此替换为你实际的 adata 调用代码)
# =====================================================================
def fetch_data_from_adata(stock_code):
# 🚨 注意:这里是模拟数据,实际请使用你的真实 adata 接口替换
end_date = datetime.today().strftime('%Y%m%d')
start_date = (datetime.today() - timedelta(days=100)).strftime('%Y%m%d')
daily_df = get_history_k(str(stock_code), start_date=start_date, end_date=end_date)
min_df = get_5min_k(stock_code, start_date='20260703', end_date='20260703')
return daily_df, min_df
# =====================================================================
# 补全缺失的类 2: 针对 Telegram 消息推送优化的子类状态机
# =====================================================================
class TelegramQuantStateMachine(QuantStateMachine):
def generate_telegram_report(self, daily_df, min_df=None):
from quant_strategy import FeatureEngine, StateRuleBook
engine = FeatureEngine(daily_df, min_df)
f = engine.calculate_all_features()
next_state, comment = StateRuleBook.evaluate_next_state(self.current_state, f)
self.current_state = next_state
msg = f"📊 *【量化做T复盘报告】* \n"
msg += f"🤖 股票代码: `{self.stock_code}`\n"
msg += f"🕒 诊断时间: {datetime.now().strftime('%Y-%m-%d %H:%M')}\n"
msg += f"📈 因子: `%B`={f['percent_b']:.2f} | 分时偏离={f['min_bias']:.2%} | 带宽分位={f['bw_quantile']:.2%}\n"
msg += f"🔍 状态: *{self.current_state}*\n"
msg += f"📝 诊断: _{comment}_\n"
msg += "-" * 30 + "\n"
msg += f"💡 *[明日手动操作指南]*:\n"
if self.current_state in ["STATE_1_OSCILLATION", "STATE_1_OSCILLATION_SQUEEZE"]:
if self.current_state == "STATE_1_OSCILLATION_SQUEEZE":
msg += "⚠️ *[变盘警告]*:弹簧已压紧,随时大突破,做T手速要快!\n"
if f['percent_b'] >= 1.0 or f['min_bias'] > 0.035:
msg += f"🟢 *【建议高抛】*:当前处于震荡高位(价格:{f['close']}),建议尾盘或明日开盘*手动卖出网格仓*!"
elif f['percent_b'] <= 0.0 or f['min_bias'] < -0.035:
msg += f"🔴 *【建议低吸】*:当前处于震荡超跌区(价格:{f['close']}),建议尾盘或明日开盘*手动买回筹码*!"
else:
msg += "⚪ *【建议观望】*:处于安全中枢内,未触及边界,明天*不要乱动*。"
elif self.current_state == "STATE_3_MAIN_WAVE":
msg += "🔥 *【强烈建议死守】*:科技股主升浪狂飙中!*禁止日内做T高抛*,锁仓死拿,享受主升浪最大利润!"
elif self.current_state == "STATE_4_WAVE_END":
msg += "⚠️ *【建议减仓】*:主升浪确认破位结束。建议手动*大举高抛/清空做T仓位*,落袋为安。"
elif self.current_state == "STATE_2_DOWN_BREAK":
msg += "🛑 *【严禁抄底】*:技术形态向下破位崩塌!明天*千万不要低吸接飞刀*,保持观望。"
elif self.current_state == "STATE_5_BOTTOMING":
msg += "🌱 *【建议潜伏】*:个股地量筑底阶段。不建议激进日内做T,但适合长线资金手动分批*定投埋伏*。"
return msg
# =====================================================================
# 补全缺失的函数 3: 数据库状态加载与固化
# =====================================================================
def load_saved_state(stock_code):
conn = sqlite3.connect(DB_FILE)
cursor = conn.cursor()
cursor.execute("SELECT current_state FROM stock_states WHERE stock_code = ?", (stock_code,))
row = cursor.fetchone()
conn.close()
return row[0] if row else "STATE_1_OSCILLATION"
def save_current_state(stock_code, state):
conn = sqlite3.connect(DB_FILE)
cursor = conn.cursor()
cursor.execute("INSERT OR REPLACE INTO stock_states (stock_code, current_state, update_time) VALUES (?, ?, ?)",
(stock_code, state, datetime.now().strftime("%Y-%m-%d %H:%M:%S")))
conn.commit()
conn.close()
# =====================================================================
# 主运行入口
# =====================================================================
def main():
logging.info("量化定时任务触发...")
try:
conn = sqlite3.connect(DB_FILE)
df = pd.read_sql_query("SELECT stock_code FROM watchlist WHERE is_active = 1", conn)
conn.close()
watchlist = df['stock_code'].tolist()
except Exception as e:
logging.error(f"读取数据库自选列表失败: {e}")
return
if not watchlist:
logging.warning("当前没有激活的关注股票。")
return
for stock_code in watchlist:
try:
daily_df, min_df = fetch_data_from_adata(stock_code)
if daily_df.empty:
logging.warning(f"[{stock_code}] 数据为空,跳过")
continue
saved_state = load_saved_state(stock_code)
machine = TelegramQuantStateMachine(initial_state=saved_state)
machine.stock_code = stock_code
tg_report = machine.generate_telegram_report(daily_df, min_df)
# 推送大字报到 TG
url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendMessage"
proxies = None
if PROXY_ENABLED and PROXY_URL:
proxies = {
"http": PROXY_URL,
"https": PROXY_URL
}
# 将 proxies 字典作为参数传入 requests
response = requests.post(
url,
json={"chat_id": CHAT_ID, "text": tg_report, "parse_mode": "Markdown"},
proxies=proxies,
timeout=60
)
save_current_state(stock_code, machine.current_state)
except Exception as e:
logging.error(f"[{stock_code}] 运行时异常: {e}", exc_info=True)
def init_database_safely():
conn = sqlite3.connect(DB_FILE)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS watchlist (
stock_code TEXT PRIMARY KEY,
stock_name TEXT,
is_active INTEGER DEFAULT 1
)
""")
cursor.execute("""
CREATE TABLE IF NOT EXISTS stock_states (
stock_code TEXT PRIMARY KEY,
current_state TEXT,
update_time TEXT
)
""")
conn.commit()
conn.close()
if __name__ == "__main__":
init_database_safely() # 👈 核心:让 main.py 每次运行时也自己检查并建表
main()