/
fgtmenow
/
tb
Обзор
Документация
Войти
/
fgtmenow
/
tb
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
dev
bot/objects/EventProcessor.py
214 строк
9 KB
Azarev Artem
start implement bot
15 апр 2025, 16:52
15 апр 2025, 16:52
1ac3c18
Код
Авторство
О чём код?
import concurrent.futures import queue import threading import uuid from collections import defaultdict from typing import Callable, Dict, Tuple, Optional, List from tqdm import tqdm from _common.domain.enums.Priority import Priority from _common.domain.util.SignletonMeta import SingletonMeta from _common.domain.util.synronized import synchronized from bot.actions.base.BaseAction import BaseAction class StopSignal: """ Сигнальный класс для указания сигнала остановки рабочих потоков. """ pass class EventProcessor(metaclass=SingletonMeta): """ Процессор событий для асинхронной обработки торговых действий. Реализует паттерн Singleton для централизованной обработки событий. Управляет очередями событий для каждого торгового символа и обеспечивает их параллельную обработку с помощью пула потоков. Атрибуты: _event_queues (Dict[str, queue.PriorityQueue]): Очереди событий для каждого символа _locks (Dict[str, threading.Lock]): Блокировки для синхронизации доступа _executor (ThreadPoolExecutor): Пул потоков для выполнения задач _futures (Dict[str, Future]): Футуры для отслеживания выполнения задач _stop_event (Event): Событие для сигнализации остановки всех потоков _callbacks (List[Callable]): Список колбэков для уведомления о завершении задач _total_events (int): Общее количество обработанных событий _progress_bar (tqdm): Общий прогресс-бар _is_running (bool): Флаг, указывающий, запущен ли процессор """ _instance = None _lock = threading.Lock() def __init__(self, max_workers: int = 10, empty_queue_timeout: int = 3): """ Инициализация процессора событий. Аргументы: max_workers: Максимальное количество рабочих потоков empty_queue_timeout: Таймаут (в секундах) для проверки пустых очередей """ self._event_queues: Dict[str, queue.PriorityQueue] = defaultdict(queue.PriorityQueue) self._locks: Dict[str, threading.Lock] = defaultdict(threading.Lock) self._executor = concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) self._futures: Dict[str, concurrent.futures.Future] = {} self._stop_event = threading.Event() self._callbacks: List[Callable[[str], None]] = [] self._total_events = 0 self._progress_bar = tqdm(total=0, desc="EventProcessor", position=0, leave=False, dynamic_ncols=True) self._is_running = False self._empty_queue_callback: Optional[Callable[[], None]] = None self._empty_queue_timeout = empty_queue_timeout self._empty_queue_monitor_event = threading.Event() self._empty_queue_monitor_thread = threading.Thread(target=self._monitor_empty_queues) self._sequence_number = 0 def register_callback(self, callback: Callable[[str], None]) -> None: """ Регистрация колбэка для вызова при завершении задачи. Аргументы: callback: Функция обратного вызова, принимающая ID задачи """ self._callbacks.append(callback) def unregister_callback(self, callback: Callable[[str], None]) -> None: """ Отмена регистрации ранее зарегистрированного колбэка. Аргументы: callback: Функция обратного вызова для удаления """ self._callbacks.remove(callback) def register_empty_queue_callback(self, callback: Callable[[], None]) -> None: """ Регистрация колбэка для вызова, если очередь событий пуста в течение определенного времени. Аргументы: callback: Функция обратного вызова без параметров """ self._empty_queue_callback = callback def _monitor_empty_queues(self): """ Мониторинг очередей событий и вызов колбэка, если они пусты в течение заданного таймаута. """ while not self._empty_queue_monitor_event.is_set(): all_empty = True for q in self._event_queues.values(): if not q.empty(): all_empty = False break if all_empty and self._empty_queue_callback: threading.Timer(self._empty_queue_timeout, self._empty_queue_callback).start() self._empty_queue_monitor_event.wait(timeout=self._empty_queue_timeout) def _process_event(self, task_info: Tuple[int, int, str, BaseAction]): """ Функция для обработки события. Аргументы: task_info: Кортеж с информацией о задаче (приоритет, порядковый номер, ID, данные события) """ _, _, task_id, event_data = task_info if issubclass(event_data.__class__, BaseAction): event_data.run(self) self._progress_bar.update(1) # Call all registered callbacks with the task ID for callback in self._callbacks: callback(task_id) def _worker_thread_func(self, symbol: str): """ Рабочий поток для обработки событий определенной валютной пары. Аргументы: symbol: Символ торговой пары """ while not self._stop_event.is_set(): try: task = self._event_queues[symbol].get(timeout=1) if isinstance(task[3], StopSignal): break self._process_event(task) self._event_queues[symbol].task_done() except queue.Empty: continue @synchronized def _start_worker_if_needed(self, symbol: str): """ Запускает рабочий поток для символа, если он еще не запущен. Аргументы: symbol: Символ торговой пары """ if symbol not in self._futures: future = self._executor.submit(self._worker_thread_func, symbol) self._futures[symbol] = future def stop(self): """ Останавливает всех воркеров и освобождает ресурсы. """ self._stop_event.set() self._empty_queue_monitor_event.set() for q in self._event_queues.values(): q.put((float('inf'), self._sequence_number, str(uuid.uuid4()), StopSignal())) self._sequence_number += 1 for future in self._futures.values(): future.result() self._executor.shutdown(wait=True) self._empty_queue_monitor_thread.join() self._progress_bar.close() print("\r", end="") # Очищаем строку перед выводом print("Все воркеры остановлены.") def enqueue(self, key: str, action: BaseAction, priority: Priority = Priority.USUAL ) -> str: """ Добавляет событие в очередь с указанным приоритетом. Аргументы: key: Ключ (обычно символ торговой пары) action: Действие для выполнения priority: Приоритет события (по умолчанию USUAL) Возвращает: str: UUID добавленного события """ with self._lock: event_uuid = str(uuid.uuid4()) self._event_queues[key].put((priority, self._sequence_number, event_uuid, action), block=True) self._sequence_number += 1 self._total_events += 1 self._progress_bar.total = self._total_events self._progress_bar.refresh() return event_uuid def start(self): """ Запускает обработку всех событий в очереди. """ if not self._is_running: self._is_running = True self._empty_queue_monitor_thread.start() self._start_all_workers() def _start_all_workers(self): """ Запускает воркеры для всех символов, у которых есть события в очереди. """ for symbol in self._event_queues.keys(): if not self._event_queues[symbol].empty(): self._start_worker_if_needed(symbol)