Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- import asyncio
- from concurrent.futures import ThreadPoolExecutor
- from time import sleep
- from typing import Callable, List
- import structlog
- import setup_logger
- logger = structlog.get_logger(__name__)
- class Emitter:
- 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 MySignal:
- def __init__(self):
- self._emitter = Emitter()
- def connect(self, handler: Callable):
- self._emitter.connect(handler)
- def emit(self, *args, **kwargs):
- self._emitter.emit(*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", 3), ("task2", 2), ("task3", 4), ("task4", 3), ("task5", 3),
- ("exit", 10)], start=1
- ):
- if data[0] == "exit":
- await asyncio.sleep(data[1])
- continue
- await asyncio.sleep(1)
- yield i, data
- def output_handler(name: str):
- logger.debug("<=== Work completed", name=name)
- async def main():
- with ThreadPoolExecutor(max_workers=3) as executor:
- async for i, data in async_generator():
- logger.debug("===> New task", name=data[0], timeout=data[1])
- handler = Handler()
- handler.signal.connect(output_handler)
- executor.submit(handler.worker, data[0], data[1])
- if __name__ == '__main__':
- asyncio.run(main())
Advertisement
Add Comment
Please, Sign In to add comment