-
Notifications
You must be signed in to change notification settings - Fork 35
Expand file tree
/
Copy pathstore.py
More file actions
184 lines (155 loc) · 6.5 KB
/
Copy pathstore.py
File metadata and controls
184 lines (155 loc) · 6.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
"""Persistent storage helpers built on Lance datasets."""
from __future__ import annotations
import logging
from importlib import import_module
from pathlib import Path
from typing import TYPE_CHECKING, Dict, Iterable, Mapping, Optional
import fsspec
import pyarrow as pa
if TYPE_CHECKING:
from types import ModuleType
from .config import KnowledgeGraphConfig
LOGGER = logging.getLogger(__name__)
class LanceGraphStore:
"""Manage Lance-backed tables that feed the query engine."""
def __init__(self, config: "KnowledgeGraphConfig"):
self._config = config
self._root: Path | str = config.storage_path
self._lance: Optional[ModuleType] = None
self._lance_attempted = False
# Initialize filesystem interface
# We convert to string to ensure compatibility with fsspec, but we'll
# use self._root (the original type) when reconstructing return values.
try:
self._fs, self._fs_path = fsspec.core.url_to_fs(
str(self._root), **(self.config.storage_options or {})
)
except ImportError:
# Re-raise explicit ImportError if protocol driver (e.g. gcsfs, s3fs)
# is missing
raise
@property
def config(self) -> "KnowledgeGraphConfig":
"""Return the configuration backing this store."""
return self._config
@property
def root(self) -> Path | str:
"""Return the root path for persisted datasets."""
return self._root
def ensure_layout(self) -> None:
"""Create the storage layout if it does not already exist."""
try:
self._fs.makedirs(self._fs_path, exist_ok=True)
except Exception:
# S3/GCS might not support directory creation or it might be implicit.
# We treat failure here as non-fatal if the path is actually accessible
# later,
# but usually makedirs is safe on object stores (no-op).
pass
def list_datasets(self) -> Dict[str, Path | str]:
"""Enumerate known Lance datasets."""
datasets: Dict[str, Path | str] = {}
try:
if not self._fs.exists(self._fs_path):
return datasets
infos = self._fs.ls(self._fs_path, detail=True)
except Exception as e:
# We want to swallow "not found" errors but raise others (like Auth errors)
if isinstance(e, FileNotFoundError):
return datasets
msg = str(e).lower()
if "not found" in msg or "no such file" in msg or "does not exist" in msg:
return datasets
raise
root_str = str(self._root)
for info in infos:
name = info["name"].rstrip("/")
base_name = name.split("/")[-1]
if info["type"] == "directory" and base_name.endswith(".lance"):
dataset_name = base_name[:-6]
full_path = f"{root_str.rstrip('/')}/{base_name}"
if isinstance(self._root, Path):
datasets[dataset_name] = Path(full_path)
else:
datasets[dataset_name] = full_path
return datasets
def _dataset_path(self, name: str) -> Path | str:
"""Create the canonical path for a dataset."""
safe_name = name.replace("/", "_")
if isinstance(self._root, Path):
return self._root / f"{safe_name}.lance"
return f"{self._root.rstrip('/')}/{safe_name}.lance"
def _get_lance(self) -> ModuleType:
if not self._lance_attempted:
self._lance_attempted = True
try:
module = import_module("lance")
except ImportError as e:
raise ImportError(
"Lance module is required but not installed. "
"Install it with: pip install pylance"
) from e
has_loader = hasattr(module, "dataset")
if not (has_loader):
raise ImportError(
"Installed `lance` package is missing required dataset APIs."
)
self._lance = module
if self._lance is None:
raise ImportError("Lance module failed to load")
return self._lance
def _path_exists(self, path: Path | str) -> bool:
if isinstance(path, Path):
return path.exists()
try:
fs, p = fsspec.core.url_to_fs(path)
except Exception:
# If we cannot resolve the filesystem (e.g. missing gcsfs), we should raise
# rather than assuming the path does not exist.
raise
try:
return fs.exists(p)
except Exception:
return False
def load_tables(
self,
names: Optional[Iterable[str]] = None,
) -> Mapping[str, "pa.Table"]:
"""Load Lance datasets as PyArrow tables.
When specific names are provided, this method computes paths directly
without enumerating all datasets - significantly faster on cloud storage.
"""
lance = self._get_lance()
self.ensure_layout()
# Only enumerate datasets when no specific names are requested
if names is not None:
requested = list(names)
else:
available = self.list_datasets()
requested = list(available.keys())
tables: Dict[str, "pa.Table"] = {}
for name in requested:
# Compute path directly - no need to look up from enumeration
path = self._dataset_path(name)
if not self._path_exists(path):
raise FileNotFoundError(f"Dataset '{name}' not found at {path}")
dataset = lance.dataset(
str(path), storage_options=self.config.storage_options
)
table = dataset.scanner().to_table()
tables[name] = table
return tables
def write_tables(self, tables: Mapping[str, "pa.Table"]) -> None:
"""Persist PyArrow tables as Lance datasets."""
lance = self._get_lance()
self.ensure_layout()
for name, table in tables.items():
if not isinstance(table, pa.Table):
raise TypeError(
f"Dataset '{name}' must be a pyarrow.Table (got {type(table)!r})"
)
path = self._dataset_path(name)
mode = "overwrite" if self._path_exists(path) else "create"
lance.write_dataset(
table, str(path), mode=mode, storage_options=self.config.storage_options
)