Writer of the Queue buffer.
| 9 | |
| 10 | |
| 11 | class QueueWriter(BufferWriter): |
| 12 | """Writer of the Queue buffer.""" |
| 13 | |
| 14 | def __init__(self, config: StorageConfig): |
| 15 | assert config.storage_type == StorageType.QUEUE.value |
| 16 | self.queue = QueueStorage.get_wrapper(config) |
| 17 | |
| 18 | async def write(self, data: List[Experience]) -> None: |
| 19 | return await self.queue.put_batch.remote(Experience.serialize_many(data)) |
| 20 | |
| 21 | async def acquire(self) -> int: |
| 22 | return await self.queue.acquire.remote() |
| 23 | |
| 24 | async def release(self) -> int: |
| 25 | return await self.queue.release.remote() |
no outgoing calls