Alpaca One-Shot 异步价格监控系统
本卡片记录专为日内标的(如 QQQ 等美股/ETF)价格监控设计的高可用、异步(Async/Await)One-Shot 自动销毁监控系统的 app.py 源码、环境参数、部署命令及 API 调用限制。
一、 核心功能设计 (One-Shot & 自毁机制)
- 单次触发即停 (One-Shot):满足条件(低于跌破价或高于涨破价)后,瞬间向飞书发送警报通知,并立即将任务激活状态重置为
is_active = False退出监控,防止价格在边界线反复摩擦带来消息轰炸。 - 当日开盘有效 (单日自毁):通过美东时间(America/New_York)时区感知,检测到已收盘(常规 16:00,或含盘前盘后 20:00)或跨天变动时,自动将监控挂起并写入磁盘,收市后自动释放资源,对后续无交易日静默。
- IO 并发锁与原子写入:在修改状态写入磁盘时,引入
asyncio.Lock()协程锁,且写入过程通过临时文件写入再调用os.replace原子替换,规避并发写盘冲突和突然断电导致的数据损耗崩溃。 - 接口安全与只读设计:系统不具备任何交易、开仓、平仓函数,并配置
X-Hermes-Token进行安全头鉴权,防范未授权的第三方越权篡改。
二、 核心监控服务源码 (app.py)
import os
import json
import logging
import asyncio
from datetime import datetime
from zoneinfo import ZoneInfo
import httpx
import aiofiles
from fastapi import FastAPI, HTTPException, Depends
from fastapi.security.api_key import APIKeyHeader
from pydantic import BaseModel
from contextlib import asynccontextmanager
# ================= 1. 系统初始化与并发锁 =================
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')
logger = logging.getLogger(__name__)
STATE_FILE = "monitor_state.json"
STATE_FILE_TMP = f"{STATE_FILE}.tmp"
NY_TZ = ZoneInfo("America/New_York")
API_KEY = os.getenv("HERMES_MONITOR_TOKEN", "DefaultSecretToken123")
api_key_header = APIKeyHeader(name="X-Hermes-Token", auto_error=False)
state_lock = asyncio.Lock()
monitor_state = {
"symbol": "QQQ",
"target_price": None,
"direction": "BELOW",
"extended_hours": False,
"is_active": False,
"created_date": None
}
async def save_state_to_disk():
async with state_lock:
try:
state_json = json.dumps(monitor_state, ensure_ascii=False, indent=4)
async with aiofiles.open(STATE_FILE_TMP, "w", encoding="utf-8") as f:
await f.write(state_json)
await f.flush()
try:
os.fsync(f.fileno())
except Exception:
pass
os.replace(STATE_FILE_TMP, STATE_FILE)
except Exception as e:
logger.error(f"💾 状态持久化失败: {e}")
async def load_state_from_disk():
global monitor_state
if os.path.exists(STATE_FILE):
try:
async with state_lock:
async with aiofiles.open(STATE_FILE, "r", encoding="utf-8") as f:
content = await f.read()
if content.strip():
monitor_state.update(json.loads(content))
logger.info(f"💾 成功恢复状态: {monitor_state}")
except json.JSONDecodeError:
logger.error("💾 状态文件已损坏,将使用默认状态并尝试覆盖修复。")
except Exception as e:
logger.error(f"💾 读取状态失败: {e}")
# ================= 2. 异步网络组件 =================
APCA_API_KEY_ID = os.getenv("APCA_API_KEY_ID", "你的APCA_API_KEY_ID")
APCA_API_SECRET_KEY = os.getenv("APCA_API_SECRET_KEY", "你的APCA_API_SECRET_KEY")
PROXY_URL = os.getenv("SOCKS_PROXY", None)
HERMES_API_URL = os.getenv("HERMES_API_URL", "http://127.0.0.1:8317/api/send_message")
FEISHU_USER_ID = os.getenv("FEISHU_USER_ID", "ou_5a1339ced5fdcdfdbea4b784c6d92c9e")
client = httpx.AsyncClient(
limits=httpx.Limits(max_keepalive_connections=5, max_connections=50),
transport=httpx.AsyncHTTPTransport(retries=3),
proxy=PROXY_URL if PROXY_URL else None,
headers={
"APCA-API-KEY-ID": APCA_API_KEY_ID,
"APCA-API-SECRET-KEY": APCA_API_SECRET_KEY,
"Accept": "application/json"
}
)
async def get_market_price(symbol: str) -> float | None:
try:
url = f"https://data.alpaca.markets/v2/stocks/{symbol}/trades/latest"
resp = await client.get(url, timeout=5.0)
resp.raise_for_status()
price = resp.json().get("trade", {}).get("p")
return float(price) if price is not None else None
except Exception as e:
logger.error(f"❌ 行情解析故障 (Alpaca Data API): {e}")
return None
async def send_feishu_notification(message: str):
try:
payload = {
"chat_id": FEISHU_USER_ID,
"text": message
}
await client.post(HERMES_API_URL, json=payload, timeout=5.0)
logger.info("🔔 飞书本地通知推送成功。")
except Exception as e:
logger.error(f"❌ 飞书本地通知推送失败: {e}")
# ================= 3. 核心监控逻辑 =================
def check_time_window(ny_now: datetime, allow_extended: bool):
if ny_now.weekday() >= 5: return False, "非交易日"
if monitor_state["created_date"] and monitor_state["created_date"] != ny_now.strftime("%Y-%m-%d"):
return False, "单日自毁"
open_hr, close_hr = (4, 20) if allow_extended else (9, 16)
market_open = ny_now.replace(hour=open_hr, minute=30 if not allow_extended else 0, second=0, microsecond=0)
market_close = ny_now.replace(hour=close_hr, minute=0, second=0, microsecond=0)
if ny_now < market_open: return False, "盘前等待中..."
if ny_now >= market_close: return False, "今日已收盘"
return True, "交易窗口内"
async def monitoring_loop():
logger.info("🚀 异步极简版监控引擎启动 (One-Shot 模式)...")
while True:
try:
if not monitor_state["is_active"]:
await asyncio.sleep(2)
continue
ny_now = datetime.now(NY_TZ)
is_in_window, reason = check_time_window(ny_now, monitor_state.get("extended_hours", False))
if not is_in_window:
if "单日自毁" in reason or "今日已收盘" in reason:
logger.info(f"🔒 {reason},当日监控任务结束,自动退出监控。")
monitor_state["is_active"] = False
await save_state_to_disk()
await asyncio.sleep(5)
continue
sym, target = monitor_state["symbol"], monitor_state["target_price"]
direction = monitor_state["direction"]
price = await get_market_price(sym)
if price and price > 1.0:
trigger = (direction == "BELOW" and price <= target) or (direction == "ABOVE" and price >= target)
if trigger:
dir_txt = "跌破" if direction == "BELOW" else "涨破"
msg = f"🚨【Hermes 信号提示】\n标的: {sym}\n现价: {price} 已{dir_txt}目标价 {target}。\n\n✅ 监控任务已完成并自动退出。"
await send_feishu_notification(msg)
monitor_state["is_active"] = False
logger.info("✅ 目标已触发,One-Shot 任务终止。")
await save_state_to_disk()
await asyncio.sleep(5)
except asyncio.CancelledError:
logger.info("🛑 收到取消信号,监控循环安全退出。")
raise
except Exception as e:
logger.error(f"💥 监控协程异常: {e}")
await asyncio.sleep(10)
# ================= 4. FastAPI 接口 =================
@asynccontextmanager
async def lifespan(app: FastAPI):
await load_state_from_disk()
task = asyncio.create_task(monitoring_loop())
def handle_task_result(t):
try:
t.result()
except asyncio.CancelledError:
pass
except Exception as err:
logger.critical(f"💥 后台监控任务异常崩溃: {err}", exc_info=True)
task.add_done_callback(handle_task_result)
yield
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
await client.aclose()
app = FastAPI(title="Hermes Agent API - OneShot", lifespan=lifespan)
class MonitorTask(BaseModel):
symbol: str = "QQQ"
target_price: float
direction: str = "BELOW"
extended_hours: bool = False
def verify_token(header_token: str = Depends(api_key_header)):
if header_token != API_KEY: raise HTTPException(status_code=403, detail="Forbidden")
return header_token
@app.post("/set_task", dependencies=[Depends(verify_token)])
async def set_task(task: MonitorTask):
if task.direction.upper() not in ["BELOW", "ABOVE"]:
raise HTTPException(status_code=400, detail="Direction must be BELOW or ABOVE")
monitor_state.update({
"symbol": task.symbol.upper(),
"target_price": task.target_price,
"direction": task.direction.upper(),
"extended_hours": task.extended_hours,
"is_active": True,
"created_date": datetime.now(NY_TZ).strftime("%Y-%m-%d")
})
await save_state_to_disk()
logger.info(f"📥 Agent 下发任务 -> {task.symbol} {task.direction} {task.target_price}")
return {"status": "success", "data": monitor_state}
@app.post("/stop_task", dependencies=[Depends(verify_token)])
async def stop_task():
monitor_state["is_active"] = False
await save_state_to_disk()
return {"status": "success", "message": "Monitoring stopped"}
@app.get("/status", dependencies=[Depends(verify_token)])
async def get_status():
return monitor_state
if __name__ == "__main__":
import uvicorn
uvicorn.run("app:app", host="127.0.0.1", port=8000, reload=False, workers=1)三、 Alpaca API 限流与数据量说明
- 频率上限 (Rate Limit):
- 免费版账户:限制为 每分钟 200 次 API 请求。
- 付费版账户:限制提升至 每分钟 1000 - 10000 次。
- 本脚本消耗计算:
- 采用 5 秒轮询间隔。每分钟发送 12 次请求(仅占用免费版限制的 6%),安全空间充裕,绝不影响并行交易程序的 API 额度。
- 免费版数据源局限:
- 免费版的实时数据源仅包含 IEX 交易所 的成交,当成交发生在其他主流交易所时,数据获取可能产生几秒的细微滞后或些许价差,但不妨碍大波动的宏观监控。
四、 宿主机 Systemd 守护运行配置
1. 服务配置文件 /etc/systemd/system/hermes-monitor.service
[Unit]
Description=Hermes Agent One-Shot Price Monitor Service
After=network.target
[Service]
Type=simple
User=root
WorkingDirectory=/root/HermesAgent/price-monitor
Environment="HERMES_MONITOR_TOKEN=DefaultSecretToken123"
Environment="APCA_API_KEY_ID=[REDACTED]"
Environment="APCA_API_SECRET_KEY=[REDACTED]"
Environment="HERMES_API_URL=http://127.0.0.1:8317/api/send_message"
Environment="FEISHU_USER_ID=[REDACTED]"
ExecStart=/usr/bin/python3 app.py
Restart=always
RestartSec=5
[Install]
WantedBy=multi-user.target2. 状态控制常用命令
# 重载系统守护进程并启动服务
systemctl daemon-reload
systemctl enable hermes-monitor --now
# 状态与日志监控
systemctl status hermes-monitor
journalctl -u hermes-monitor -n 50 --no-pager相关链接
- Streamlit多因子回测框架启动配置 Tushare 数据源限流防御及本地回测启动
- 量化机器人审计与缺陷报告 隐波及 Delta 本地容灾计算机制
- DockerCompose工作流与更新 常用运维指令