Get a buffer writer for the given dataset name.
(config: BufferStorageConfig)
| 24 | |
| 25 | |
| 26 | def get_buffer_writer(config: BufferStorageConfig) -> BufferWriter: |
| 27 | """Get a buffer writer for the given dataset name.""" |
| 28 | if not isinstance(config, StorageConfig): |
| 29 | storage_config: StorageConfig = config.to_storage_config() |
| 30 | else: |
| 31 | storage_config = config |
| 32 | if storage_config.storage_type == StorageType.SQL.value: |
| 33 | from trinity.buffer.writer.sql_writer import SQLWriter |
| 34 | |
| 35 | return SQLWriter(storage_config) |
| 36 | elif storage_config.storage_type == StorageType.QUEUE.value: |
| 37 | from trinity.buffer.writer.queue_writer import QueueWriter |
| 38 | |
| 39 | return QueueWriter(storage_config) |
| 40 | elif storage_config.storage_type == StorageType.FILE.value: |
| 41 | from trinity.buffer.writer.file_writer import JSONWriter |
| 42 | |
| 43 | return JSONWriter(storage_config) |
| 44 | else: |
| 45 | raise ValueError(f"{storage_config.storage_type} not supported.") |