import asyncio import csv import logging from pathlib import Path from typing import Any, Collection from .base import BaseWriter logger = logging.getLogger('users_cleanup') CsvDictWriterT = dict[str, Any] class StringWriter(BaseWriter[str]): def _write_record(self, record: str) -> None: assert self._fp is not None, 'File is not opened. Probably you forgot to use context manager.' self._fp.write(record) class CsvDictWriter(BaseWriter[CsvDictWriterT]): def __init__(self, *, path: Path, queue: asyncio.Queue[CsvDictWriterT], header_keys: Collection[str]) -> None: super().__init__(path=path, queue=queue) self.header_keys = header_keys self.writer: csv.DictWriter[Any] | None = None def __enter__(self) -> 'CsvDictWriter': super().__enter__() assert self._fp is not None, 'File is not opened. Probably you forgot to use context manager.' self.writer = csv.DictWriter(self._fp, self.header_keys) self.writer.writeheader() return self def _write_record(self, record: CsvDictWriterT) -> None: assert self.writer is not None self.writer.writerow(record) async def write(self, path: Path, queue: asyncio.Queue[CsvDictWriterT]) -> None: logger.info(f'Writing results to: {path}') path.parent.mkdir(parents=True, exist_ok=True) with path.open('w') as f: writer = csv.DictWriter(f, fieldnames=self.header_keys) writer.writeheader() try: while True: record = await queue.get() writer.writerow(record) queue.task_done() except asyncio.QueueShutDown: logger.debug('Result writer finished.')