Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- import asyncio
- import threading
- from enum import StrEnum
- from time import sleep
- from typing import Callable, List, Dict, Optional, Union, Any
- import structlog
- import setup_logger
- logger = structlog.get_logger(__name__)
- class MySignal:
- def __init__(self):
- self.handlers: List[Callable] = []
- def connect(self, handler: Callable):
- self.handlers.append(handler)
- def emit(self, *args, **kwargs):
- for handler in self.handlers:
- handler(*args, **kwargs)
- class Handler:
- def __init__(self):
- self.signal = MySignal()
- def worker(self, name: str, timeout: int):
- sleep(timeout)
- self.signal.emit(name)
- async def async_generator():
- # Имитация NATS Subscribe
- for i, data in enumerate([
- ("task1", 10), ("task2", 7), ("status", "task2"), ("status", "no_task"), ("task3", 6), ("status", "task3"),
- ("task4", 3), ("status", "task2"), ("task5", 3),
- ("exit", 1)], start=1
- ):
- if data[0] == "exit":
- await asyncio.sleep(data[1])
- continue
- await asyncio.sleep(1)
- yield i, data
- class WorkerThread(threading.Thread):
- def __init__(self, queue: 'Queue'):
- threading.Thread.__init__(self)
- self._exit = False
- self.queue = queue
- self.start()
- def run(self):
- while not self._exit:
- task: Optional[tuple[Any, ...]] = self.queue.get_task()
- if task is not None:
- name, data = task
- self.queue.set_work(name, self)
- handler, args, kwargs = data
- handler(*args, **kwargs)
- self.queue.done_work(name)
- class Status(StrEnum):
- RUNNING = "running"
- WAITING = "waiting"
- DONE = "done"
- class Queue:
- def __init__(self, max_thread: int = 0):
- self._queue = {}
- self.pool = {}
- self._workers = []
- self.init_threads(max_thread)
- def __enter__(self):
- # Инициализация ресурсов или выполнение действий при входе в блок with
- logger.debug("Run queue")
- print()
- return self
- def __exit__(self, exc_type, exc_value, traceback):
- # Завершение работы, освобождение ресурсов или выполнение действий при выходе из блока with
- logger.debug("Exit...")
- # while not len(self._queue) == 0: ...
- for worker in self._workers:
- worker._exit = True
- worker.join()
- logger.debug("Exit queue")
- # Здесь вы можете выполнить необходимые действия при завершении работы класса
- def init_threads(self, max_threads: int):
- for _ in range(max_threads):
- self._workers.append(WorkerThread(self))
- logger.debug(f"Create {max_threads} threads")
- def get_task(self) -> Optional[tuple[Any, ...]]:
- try:
- return self._queue.popitem()
- except KeyError:
- return None
- def set_work(self, name: str, worker: 'WorkerThread'):
- self.pool[name] = worker
- def done_work(self, name: str):
- self.pool.pop(name, None)
- def get_status(self, name: str):
- if self._queue.get(name, None) is not None:
- return Status.WAITING
- elif self.pool.get(name, None) is not None:
- return Status.RUNNING
- return Status.DONE
- def create_task(self, __handler: Callable, __name: str, *args, **kwargs):
- self._queue[__name] = (__handler, args, kwargs)
- def output_handler(name: str):
- logger.debug("<=== Work completed", name=name)
- async def main():
- with Queue(2) as queue:
- async for i, data in async_generator():
- if data[0] == "status":
- name = data[1]
- status: Status = queue.get_status(name)
- logger.debug(f"Status {name}: {status}")
- else:
- logger.debug("===> New task", name=data[0], timeout=data[1])
- handler = Handler()
- handler.signal.connect(output_handler)
- queue.create_task(handler.worker, data[0], name=data[0], timeout=data[1])
- if __name__ == '__main__':
- asyncio.run(main())
- LOGS:
- 2024-01-24 14:28:54 [debug ] Create 2 threads [test_thread_3.py:init_threads:112]
- 2024-01-24 14:28:57 [debug ] Run queue [test_thread_3.py:__enter__:95]
- 2024-01-24 14:28:58 [debug ] ===> New task [test_thread_3.py:main:149] name=task1 timeout=10
- 2024-01-24 14:29:00 [debug ] ===> New task [test_thread_3.py:main:149] name=task2 timeout=7
- 2024-01-24 14:29:01 [debug ] Status task2: running [test_thread_3.py:main:147]
- 2024-01-24 14:29:02 [debug ] Status no_task: done [test_thread_3.py:main:147]
- 2024-01-24 14:29:03 [debug ] ===> New task [test_thread_3.py:main:149] name=task3 timeout=6
- 2024-01-24 14:29:04 [debug ] Status task3: waiting [test_thread_3.py:main:147]
- 2024-01-24 14:29:05 [debug ] ===> New task [test_thread_3.py:main:149] name=task4 timeout=3
- 2024-01-24 14:29:06 [debug ] Status task2: running [test_thread_3.py:main:147]
- 2024-01-24 14:29:07 [debug ] <=== Work completed [test_thread_3.py:output_handler:138] name=task2
- 2024-01-24 14:29:07 [debug ] ===> New task [test_thread_3.py:main:149] name=task5 timeout=3
- 2024-01-24 14:29:08 [debug ] Exit... [test_thread_3.py:__exit__:101]
- 2024-01-24 14:29:09 [debug ] <=== Work completed [test_thread_3.py:output_handler:138] name=task1
- 2024-01-24 14:29:10 [debug ] <=== Work completed [test_thread_3.py:output_handler:138] name=task4
- 2024-01-24 14:29:12 [debug ] <=== Work completed [test_thread_3.py:output_handler:138] name=task5
- 2024-01-24 14:29:12 [debug ] Exit queue [test_thread_3.py:__exit__:106]
Advertisement
Add Comment
Please, Sign In to add comment