|
| 1 | +import asyncio |
| 2 | +import ydb.aio |
| 3 | +import time |
| 4 | +import logging |
| 5 | +from aiolimiter import AsyncLimiter |
| 6 | + |
| 7 | +from .base import BaseJobManager |
| 8 | +from core.metrics import OP_TYPE_READ, OP_TYPE_WRITE |
| 9 | + |
| 10 | +logger = logging.getLogger(__name__) |
| 11 | + |
| 12 | + |
| 13 | +class AsyncTopicJobManager(BaseJobManager): |
| 14 | + def __init__(self, driver, args, metrics): |
| 15 | + super().__init__(driver, args, metrics) |
| 16 | + self.driver: ydb.aio.Driver = driver |
| 17 | + |
| 18 | + async def run_tests(self): |
| 19 | + tasks = [ |
| 20 | + *await self._run_topic_write_jobs(), |
| 21 | + *await self._run_topic_read_jobs(), |
| 22 | + *self._run_metric_job(), |
| 23 | + ] |
| 24 | + |
| 25 | + await asyncio.gather(*tasks) |
| 26 | + |
| 27 | + async def _run_topic_write_jobs(self): |
| 28 | + logger.info("Start async topic write jobs") |
| 29 | + |
| 30 | + write_limiter = AsyncLimiter(max_rate=self.args.write_rps, time_period=1) |
| 31 | + |
| 32 | + tasks = [] |
| 33 | + for i in range(self.args.write_threads): |
| 34 | + task = asyncio.create_task(self._run_topic_writes(write_limiter, i), name=f"slo_topic_write_{i}") |
| 35 | + tasks.append(task) |
| 36 | + |
| 37 | + return tasks |
| 38 | + |
| 39 | + async def _run_topic_read_jobs(self): |
| 40 | + logger.info("Start async topic read jobs") |
| 41 | + |
| 42 | + read_limiter = AsyncLimiter(max_rate=self.args.read_rps, time_period=1) |
| 43 | + |
| 44 | + tasks = [] |
| 45 | + for i in range(self.args.read_threads): |
| 46 | + task = asyncio.create_task(self._run_topic_reads(read_limiter), name=f"slo_topic_read_{i}") |
| 47 | + tasks.append(task) |
| 48 | + |
| 49 | + return tasks |
| 50 | + |
| 51 | + async def _run_topic_writes(self, limiter, partition_id=None): |
| 52 | + start_time = time.time() |
| 53 | + logger.info("Start async topic write workload") |
| 54 | + |
| 55 | + async with self.driver.topic_client.writer( |
| 56 | + self.args.path, |
| 57 | + codec=ydb.TopicCodec.GZIP, |
| 58 | + partition_id=partition_id, |
| 59 | + ) as writer: |
| 60 | + logger.info("Async topic writer created") |
| 61 | + |
| 62 | + message_count = 0 |
| 63 | + while time.time() - start_time < self.args.time: |
| 64 | + async with limiter: |
| 65 | + message_count += 1 |
| 66 | + |
| 67 | + content = f"message_{message_count}_{asyncio.current_task().get_name()}".encode("utf-8") |
| 68 | + |
| 69 | + if len(content) < self.args.message_size: |
| 70 | + content += b"x" * (self.args.message_size - len(content)) |
| 71 | + |
| 72 | + message = ydb.TopicWriterMessage(data=content) |
| 73 | + |
| 74 | + ts = self.metrics.start((OP_TYPE_WRITE,)) |
| 75 | + try: |
| 76 | + await writer.write_with_ack(message) |
| 77 | + logger.info("Write message: %s", content) |
| 78 | + self.metrics.stop((OP_TYPE_WRITE,), ts) |
| 79 | + except Exception as e: |
| 80 | + self.metrics.stop((OP_TYPE_WRITE,), ts, error=e) |
| 81 | + logger.error("Write error: %s", e) |
| 82 | + |
| 83 | + logger.info("Stop async topic write workload") |
| 84 | + |
| 85 | + async def _run_topic_reads(self, limiter): |
| 86 | + start_time = time.time() |
| 87 | + logger.info("Start async topic read workload") |
| 88 | + |
| 89 | + async with self.driver.topic_client.reader( |
| 90 | + self.args.path, |
| 91 | + self.args.consumer, |
| 92 | + ) as reader: |
| 93 | + logger.info("Async topic reader created") |
| 94 | + |
| 95 | + while time.time() - start_time < self.args.time: |
| 96 | + async with limiter: |
| 97 | + ts = self.metrics.start((OP_TYPE_READ,)) |
| 98 | + try: |
| 99 | + msg = await reader.receive_message() |
| 100 | + if msg is not None: |
| 101 | + logger.info("Read message: %s", msg.data.decode()) |
| 102 | + await reader.commit_with_ack(msg) |
| 103 | + |
| 104 | + self.metrics.stop((OP_TYPE_READ,), ts) |
| 105 | + except Exception as e: |
| 106 | + self.metrics.stop((OP_TYPE_READ,), ts, error=e) |
| 107 | + logger.error("Read error: %s", e) |
| 108 | + |
| 109 | + logger.info("Stop async topic read workload") |
| 110 | + |
| 111 | + def _run_metric_job(self): |
| 112 | + if not self.args.prom_pgw: |
| 113 | + return [] |
| 114 | + |
| 115 | + # Create async task for metrics |
| 116 | + task = asyncio.create_task(self._async_metric_sender(self.args.time), name="slo_metrics_sender") |
| 117 | + return [task] |
| 118 | + |
| 119 | + async def _async_metric_sender(self, runtime): |
| 120 | + start_time = time.time() |
| 121 | + logger.info("Start push metrics (async)") |
| 122 | + |
| 123 | + limiter = AsyncLimiter(max_rate=10**6 // self.args.report_period, time_period=1) |
| 124 | + |
| 125 | + while time.time() - start_time < runtime: |
| 126 | + async with limiter: |
| 127 | + # Call sync metrics.push() in executor to avoid blocking |
| 128 | + await asyncio.get_event_loop().run_in_executor(None, self.metrics.push) |
| 129 | + |
| 130 | + logger.info("Stop push metrics (async)") |
0 commit comments