Download tcod/trinity/buffer/writer/queue_writer.py from SeanWang0027/ftb-sciworld-repro: direct link, hf CLI and curl.
- Browser
- Download file 873 Bytes
-
https://huggingface.co/SeanWang0027/ftb-sciworld-repro/resolve/main/tcod/trinity/buffer/writer/queue_writer.py
- Command line
-
hf download hf://SeanWang0027/ftb-sciworld-repro/tcod/trinity/buffer/writer/queue_writer.py
-
curl -L -o queue_writer.py https://huggingface.co/SeanWang0027/ftb-sciworld-repro/resolve/main/tcod/trinity/buffer/writer/queue_writer.py
873 Bytes
| """Writer of the Queue buffer.""" | |
| from typing import List | |
| import ray | |
| from trinity.buffer.buffer_writer import BufferWriter | |
| from trinity.buffer.storage.queue import QueueStorage | |
| from trinity.common.config import StorageConfig | |
| from trinity.common.constants import StorageType | |
| class QueueWriter(BufferWriter): | |
| """Writer of the Queue buffer.""" | |
| def __init__(self, config: StorageConfig): | |
| assert config.storage_type == StorageType.QUEUE.value | |
| self.queue = QueueStorage.get_wrapper(config) | |
| def write(self, data: List) -> None: | |
| ray.get(self.queue.put_batch.remote(data)) | |
| async def write_async(self, data): | |
| return await self.queue.put_batch.remote(data) | |
| async def acquire(self) -> int: | |
| return await self.queue.acquire.remote() | |
| async def release(self) -> int: | |
| return await self.queue.release.remote() | |