Files

214 lines
10 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
conn = sqlite3.connect(DB_FILE)
cursor = conn.cursor()
cursor.execute("SELECT stock_name FROM watchlist WHERE is_active = 1 and stock_code = ?", (self.stock_code,))
row = cursor.fetchone()
conn.close()
name = row[0] if row else "unknow"
msg = f"📊 *【量化做T复盘报告】* \n"
msg += f"🤖 股票代码: `{self.stock_code}`\n"
msg += f"🤖 股票名称: `{name}`\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"
## 盘整阶段,三个指标会触发卖出提醒
## kdj死叉,且cci在下降中,且收盘价在5日均线下
## cci顶背离,且k线在d线之下,且收盘价在5日均线下
## cci从100以上下降到100以下,且k线在d线之下,且收盘价在5日均线下
if ((f['K'] < f['D'] and f['prev_K'] > f['prev_D']) ## kdj死叉
and (f['cci'] < f['prev_cci']) ## cci在下降中
and (f['close'] < f['ma5'])): ## 收盘价在5日均线下
msg += f"🟢 *【建议高抛】*:当前处于震荡高位(价格:{f['close']}),建议尾盘或明日开盘*手动卖出网格仓*!"
elif ((f['close'] > f['prev_close'] and f['cci'] < f['prev_cci']) ## cci顶背离
and (f['K'] < f['D']) ## k线在d线之下
and (f['close'] < f['ma5'])): ## 收盘价在5日均线下
msg += f"🟢 *【建议高抛】*:当前处于震荡高位(价格:{f['close']}),建议尾盘或明日开盘*手动卖出网格仓*!"
elif ((f['cci'] < 100 and f['prev_cci'] > 100) ## cci从100跌到100以下
and (f['K'] < f['D']) ## k线在d线之下
and (f['close'] < f['ma5'])): ## 收盘价在5日均线下
msg += f"🟢 *【建议高抛】*:当前处于震荡高位(价格:{f['close']}),建议尾盘或明日开盘*手动卖出网格仓*!"
## 盘整阶段,两个指标会触发买入提醒
## kdj金叉,且cci连续两天上涨,且收盘价在5日均线上
## cci连续两天上涨,且k线在d线上,且收盘价在5日均线上
elif ((f['cci'] > f['prev_cci'] and f['prev_cci'] > f['prev_cci_2']) ## cci连续两天上涨
and (f['K'] > f['D'] and f['prev_K'] <= f['prev_D']) ## kdj金叉
and (f['close'] > f['ma5'])): ## 收盘价站上五日线
msg += f"🔴 *【建议低吸】*:当前处于震荡超跌区(价格:{f['close']}),建议尾盘或明日开盘*手动买回筹码*!"
elif ((f['cci'] > f['prev_cci'] and f['prev_cci'] > f['prev_cci_2'] and f['prev_cci_2'] < f['prev_cci_3'])
and (f['K'] >= f['D'])
and (f['close'] > f['ma5'])):
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()