Skip to content

Commit 8adbd89

Browse files
Update pyarrow.py
1 parent 068aae5 commit 8adbd89

1 file changed

Lines changed: 22 additions & 2 deletions

File tree

‎pyiceberg/io/pyarrow.py‎

Lines changed: 22 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@
3737
import re
3838
import uuid
3939
import warnings
40+
import weakref
4041
from abc import ABC, abstractmethod
4142
from collections.abc import Callable, Iterable, Iterator
4243
from copy import copy
@@ -396,11 +397,30 @@ def to_input_file(self) -> PyArrowFile:
396397
return self
397398

398399

400+
def _fs_by_scheme_cache(file_io: PyArrowFileIO) -> Callable[[str, str | None], FileSystem]:
401+
"""Return a cached FileSystem factory that only weakly references ``file_io``.
402+
403+
Caching the bound method ``file_io._initialize_fs`` directly would make the cache hold
404+
``file_io`` while ``file_io`` holds the cache. That reference cycle keeps the FileIO and
405+
its cached filesystems (and their connection pools) alive until the cycle collector runs.
406+
"""
407+
file_io_ref = weakref.ref(file_io)
408+
409+
@lru_cache
410+
def fs_by_scheme(scheme: str, netloc: str | None = None) -> FileSystem:
411+
io = file_io_ref()
412+
if io is None:
413+
raise ReferenceError("PyArrowFileIO has already been garbage collected")
414+
return io._initialize_fs(scheme, netloc)
415+
416+
return fs_by_scheme
417+
418+
399419
class PyArrowFileIO(FileIO):
400420
fs_by_scheme: Callable[[str, str | None], FileSystem]
401421

402422
def __init__(self, properties: Properties = EMPTY_DICT):
403-
self.fs_by_scheme: Callable[[str, str | None], FileSystem] = lru_cache(self._initialize_fs)
423+
self.fs_by_scheme: Callable[[str, str | None], FileSystem] = _fs_by_scheme_cache(self)
404424
super().__init__(properties=properties)
405425

406426
@staticmethod
@@ -725,7 +745,7 @@ def __getstate__(self) -> dict[str, Any]:
725745
def __setstate__(self, state: dict[str, Any]) -> None:
726746
"""Deserialize the state into a PyArrowFileIO instance."""
727747
self.__dict__ = state
728-
self.fs_by_scheme = lru_cache(self._initialize_fs)
748+
self.fs_by_scheme = _fs_by_scheme_cache(self)
729749

730750

731751
def schema_to_pyarrow(

0 commit comments

Comments
 (0)