В высокочастотном и алгоритмическом трейдинге сетевой сбой, зависание процесса или отключение электроэнергии на стороне сервера могут привести к катастрофическим потерям. Если бот успел выставить лимитные ордера, а затем потерял связь, эти ордера остаются на бирже без контроля. Для решения этой проблемы применяется паттерн Dead Man’s Switch (DMS) — предохранитель, который автоматически ликвидирует риски, если бот перестает подавать признаки жизни.
В то время как институциональные игроки используют встроенные механизмы протоколов, о которых мы писали в статье про интеграцию Python с FIX API (QuickFIX) для HFT торговли, для розничных и полупрофессиональных систем на базе REST/WebSocket оптимальным решением является создание независимого локального или удаленного процесса-наблюдателя (Watchdog).
Исходный код
Ниже представлена промышленная реализация архитектуры DMS на Python с использованием библиотеки asyncio. Система состоит из двух компонентов: торгового бота, отправляющего легковесные UDP-хартбиты, и независимого демона-наблюдателя (Watchdog), который мониторит активность и выполняет экстренную отмену ордеров в случае зависания основного процесса.
import asyncio
import socket
import time
import logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s [%(levelname)s] %(message)s'
)
class DeadMansSwitchWatchdog:
def __init__(self, host: str = '127.0.0.1', port: int = 9999, timeout: float = 3.0):
self.host = host
self.port = port
self.timeout = timeout
self.last_heartbeat = time.time()
self.is_active = True
self.transport = None
async def start(self):
loop = asyncio.get_running_loop()
self.transport, protocol = await loop.create_datagram_endpoint(
lambda: WatchdogProtocol(self),
local_addr=(self.host, self.port)
)
logging.info(f'Watchdog запущен на {self.host}:{self.port}')
try:
await self.monitor_loop()
finally:
self.stop()
async def monitor_loop(self):
while self.is_active:
await asyncio.sleep(0.5)
elapsed = time.time() - self.last_heartbeat
if elapsed > self.timeout:
logging.error(f'ВНИМАНИЕ: Хартбит отсутствует {elapsed:.2f} сек! Запуск экстренной отмены...')
await self.trigger_emergency_action()
self.is_active = False
def feed(self):
self.last_heartbeat = time.time()
async def trigger_emergency_action(self):
# Здесь вызывается API биржи для отмены всех ордеров
# Например, отправка ордеров на Bybit V5
logging.warning('ЭКСТРЕННЫЙ ВЫЗОВ: Все активные ордера аннулированы, позиции переведены в безопасный режим.')
def stop(self):
self.is_active = False
if self.transport:
self.transport.close()
logging.info('Watchdog остановлен.')
class WatchdogProtocol(asyncio.DatagramProtocol):
def __init__(self, watchdog: DeadMansSwitchWatchdog):
self.watchdog = watchdog
def datagram_received(self, data: bytes, addr):
if data == b'PING':
self.watchdog.feed()
class TradingBot:
def __init__(self, watchdog_host: str = '127.0.0.1', watchdog_port: int = 9999, heartbeat_interval: float = 1.0):
self.watchdog_host = watchdog_host
self.watchdog_port = watchdog_port
self.heartbeat_interval = heartbeat_interval
self.is_running = True
async def run(self):
# Используем UDP-сокет для отправки неблокирующих хартбитов
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
logging.info('Торговый бот запущен. Начало отправки хартбитов...')
try:
while self.is_running:
sock.sendto(b'PING', (self.watchdog_host, self.watchdog_port))
logging.info('Хартбит отправлен на Watchdog')
# Имитация работы торговой логики
await asyncio.sleep(self.heartbeat_interval)
except asyncio.CancelledError:
logging.info('Торговый бот остановлен пользователем.')
finally:
sock.close()
async def main():
# Запуск Watchdog и Бота в едином event loop для демонстрации.
# В продакшене это должны быть два абсолютно независимых процесса!
watchdog = DeadMansSwitchWatchdog(timeout=3.0)
bot = TradingBot(heartbeat_interval=1.0)
# Запускаем задачи асинхронно
watchdog_task = asyncio.create_task(watchdog.start())
bot_task = asyncio.create_task(bot.run())
# Симулируем 5 секунд нормальной работы
await asyncio.sleep(5.0)
# Симулируем зависание/падение бота
logging.warning('Симуляция падения торгового бота...')
bot_task.cancel()
# Ожидаем срабатывания триггера Watchdog
await asyncio.sleep(4.0)
watchdog_task.cancel()
if __name__ == '__main__':
try:
asyncio.run(main())
except KeyboardInterrupt:
pass Разбор параметров
hostиport: Сетевой адрес и порт для UDP-сокета. Использование протокола UDP критически важно, так как он не требует установки соединения (handshake) и не блокирует основной поток выполнения торгового бота при сетевых задержках.timeout: Максимально допустимое время отсутствия сигнала (хартбита) от бота в секундах. Для HFT-систем этот параметр может составлять0.5–1.0секунды, для среднесрочных стратегий —3.0–5.0секунд.heartbeat_interval: Частота отправки сигналов ботом. Она должна быть как минимум в 2-3 раза меньше, чемtimeout, чтобы избежать ложных срабатываний из-за кратковременных сетевых джиттеров.trigger_emergency_action: Метод, выполняющий экстренное закрытие рисков. В реальных условиях здесь реализуется отправка запросов на биржу. Например, для минимизации проскальзываний и защиты лимитных заявок можно использовать логику, описанную в статье про создание Post-Only и Reduce-Only ордеров на Bybit V5 в Python.
Как запустить
Для тестирования механизма выполните следующие шаги:
1. Скопируйте представленный код в файл dms_protection.py.
2. Запустите скрипт с помощью интерпретатора Python:
python dms_protection.py 3. В консоли вы увидите логи нормального обмена сообщениями, после чего произойдет симуляция падения бота. Ровно через 3 секунды после последнего хартбита Watchdog зафиксирует аварию и вызовет метод экстренной очистки ордеров.
В промышленной эксплуатации рекомендуется запускать Watchdog на отдельном сервере (co-location), географически близком к бирже. Это гарантирует, что даже при полном отключении дата-центра, где запущен ваш бот, или при падении глобального интернета, внешний наблюдатель сможет оперативно очистить стакан от ваших заявок. Для анализа рыночной конъюнктуры перед перезапуском бота вы также можете использовать внешние метрики, например, настроив интеграцию Python с Glassnode API для получения SOPR и NVT.




