|
1 | 1 | import logging |
2 | 2 |
|
| 3 | +from fsspec.asyn import AbstractAsyncStreamedFile |
3 | 4 | from fsspec.spec import AbstractBufferedFile |
| 5 | +from opendal import AsyncFile as OpendalAsyncFile |
4 | 6 | from opendal import File as OpendalFile |
5 | 7 |
|
6 | 8 | logger = logging.getLogger("opendalfs") |
@@ -148,3 +150,133 @@ def close(self): |
148 | 150 | self._opendal_writer.close() |
149 | 151 | finally: |
150 | 152 | self._opendal_writer = None |
| 153 | + |
| 154 | + |
| 155 | +class OpendalAsyncBufferedFile(AbstractAsyncStreamedFile): |
| 156 | + """Async buffered file implementation for OpenDAL.""" |
| 157 | + |
| 158 | + _opendal_writer: OpendalAsyncFile | None |
| 159 | + _append_via_write: bool |
| 160 | + _initiated: bool |
| 161 | + _exclusive_create: bool |
| 162 | + |
| 163 | + def __init__( |
| 164 | + self, |
| 165 | + fs, |
| 166 | + path, |
| 167 | + mode="rb", |
| 168 | + block_size="default", |
| 169 | + autocommit=True, |
| 170 | + cache_type="readahead", |
| 171 | + cache_options=None, |
| 172 | + size=None, |
| 173 | + **kwargs, |
| 174 | + ): |
| 175 | + self._exclusive_create = mode == "xb" |
| 176 | + normalized_mode = "wb" if self._exclusive_create else mode |
| 177 | + super().__init__( |
| 178 | + fs, |
| 179 | + path, |
| 180 | + mode=normalized_mode, |
| 181 | + block_size=block_size, |
| 182 | + autocommit=autocommit, |
| 183 | + cache_type=cache_type, |
| 184 | + cache_options=cache_options, |
| 185 | + size=size, |
| 186 | + **kwargs, |
| 187 | + ) |
| 188 | + |
| 189 | + self._opendal_writer = None |
| 190 | + self._append_via_write = False |
| 191 | + self._initiated = False |
| 192 | + |
| 193 | + async def _fetch_range(self, start: int, end: int): |
| 194 | + if start >= end: |
| 195 | + return b"" |
| 196 | + |
| 197 | + reader = await self.fs.async_fs.open(self.path, "rb") |
| 198 | + try: |
| 199 | + await reader.seek(start) |
| 200 | + return await reader.read(end - start) |
| 201 | + finally: |
| 202 | + await reader.close() |
| 203 | + |
| 204 | + async def _upload_chunk(self, final: bool = False): |
| 205 | + if not self._initiated: |
| 206 | + raise RuntimeError("Upload has not been initiated") |
| 207 | + |
| 208 | + self.buffer.seek(0) |
| 209 | + chunk = self.buffer.read() |
| 210 | + |
| 211 | + if not chunk: |
| 212 | + if not final: |
| 213 | + return False |
| 214 | + if self.mode == "ab" and self._append_via_write: |
| 215 | + if not await self.fs.async_fs.exists(self.path): |
| 216 | + await self.fs.async_fs.write(self.path, b"") |
| 217 | + return None |
| 218 | + await self._commit_upload() |
| 219 | + return None |
| 220 | + |
| 221 | + if self.mode == "ab" and self._append_via_write: |
| 222 | + await self.fs.async_fs.write(self.path, chunk, append=True) |
| 223 | + return None |
| 224 | + |
| 225 | + if self._opendal_writer is None: |
| 226 | + self._opendal_writer = await self.fs.async_fs.open(self.path, "wb") |
| 227 | + |
| 228 | + await self._opendal_writer.write(chunk) |
| 229 | + |
| 230 | + if final: |
| 231 | + await self._commit_upload() |
| 232 | + return None |
| 233 | + |
| 234 | + async def _initiate_upload(self) -> None: |
| 235 | + if self._initiated: |
| 236 | + return |
| 237 | + |
| 238 | + if self._exclusive_create and await self.fs.async_fs.exists(self.path): |
| 239 | + raise FileExistsError(self.path) |
| 240 | + |
| 241 | + if self.mode == "ab": |
| 242 | + cap = self.fs.async_fs.capability() |
| 243 | + if getattr(cap, "write_can_append", False): |
| 244 | + self._append_via_write = True |
| 245 | + self.offset = self.loc |
| 246 | + else: |
| 247 | + try: |
| 248 | + existing = await self.fs.async_fs.read(self.path) |
| 249 | + except FileNotFoundError: |
| 250 | + existing = b"" |
| 251 | + if existing: |
| 252 | + self._opendal_writer = await self.fs.async_fs.open(self.path, "wb") |
| 253 | + await self._opendal_writer.write(existing) |
| 254 | + self.offset = len(existing) |
| 255 | + |
| 256 | + self._initiated = True |
| 257 | + |
| 258 | + async def _commit_upload(self) -> None: |
| 259 | + if self.mode == "ab" and self._append_via_write: |
| 260 | + return |
| 261 | + |
| 262 | + if self._opendal_writer is None: |
| 263 | + await self.fs.async_fs.write(self.path, b"") |
| 264 | + return |
| 265 | + |
| 266 | + try: |
| 267 | + await self._opendal_writer.close() |
| 268 | + finally: |
| 269 | + self._opendal_writer = None |
| 270 | + |
| 271 | + async def close(self): |
| 272 | + if self.closed: |
| 273 | + return |
| 274 | + |
| 275 | + try: |
| 276 | + await super().close() |
| 277 | + finally: |
| 278 | + if self._opendal_writer is not None: |
| 279 | + try: |
| 280 | + await self._opendal_writer.close() |
| 281 | + finally: |
| 282 | + self._opendal_writer = None |
0 commit comments