Skip to content

Commit c4ebceb

Browse files
committed
Move remote Parquet tables under RemoteStore
1 parent 0fdc62c commit c4ebceb

12 files changed

Lines changed: 1004 additions & 1181 deletions

‎doc/guides/remote_tables.md‎

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ Fixed-width, `blosc2.utf8()`, batch-backed variable-length, list, struct/object,
55
and dictionary columns are fetched on demand, including their null masks. A table
66
inside a hierarchy can also be opened through `RemoteStore`.
77

8-
## Single-file Parquet (experimental)
8+
## Single-file Parquet
99

1010
`blosc2.open()` recognizes `.parquet` paths, including fsspec URLs, and returns
1111
a read-only `RemoteCTable`. Use `source_format="parquet"` for an extensionless
@@ -32,14 +32,28 @@ opens and cached reads need no connection to the source. Cached sources are
3232
assumed unchanged until `refresh()` is called. Use `lazy=False` to import the
3333
whole table eagerly.
3434

35+
Parquet is a `RemoteStore` with one root CTable. Use
36+
`RemoteStore(url, allow_table_root=True)` and `store[""]` when a store operation
37+
needs to own the cache and traffic counters. A Parquet file accepts only the
38+
root selector (`""` or `"/"`); selecting a child path raises an error.
39+
`RemoteCTable(url)` and lazy `blosc2.open(url)` use the same owner. NONE retains
40+
no converted row groups, MEMORY shares the store budget, and DISK retains
41+
complete converted physical-column/row-group units. A small slice can therefore
42+
read a whole row group. Shared DISK caches support `read_cached_table()` and
43+
offline `trim_sparse_cache()` with the same aggregate allowance.
44+
Caches from the earlier Parquet prototype layout are not migrated; use a fresh
45+
`cache_dir` for this RemoteStore layout.
46+
3547
HTTP servers must honor byte-range requests; fsspec raises a range-request error
3648
for servers that only return complete files. Download the file and import it
3749
locally when range access is unavailable.
3850

39-
`table.save("reference.b2z")` creates a portable archive with data already in
51+
`table.save("reference.b2z")` creates a RemoteStore archive with data already in
4052
the cache. Opening it reuses that data and fetches missing groups on demand.
41-
Use `include_cache=False` to save only the table metadata. Older `.b2nd`
42-
references remain readable.
53+
Archives of local sources keep an absolute path to the local Parquet file.
54+
Use `include_cache=False` to save only the table metadata. `.b2z` exports use
55+
the RemoteStore version 1 manifest with `kind="parquet"`. Runtime filesystems, credentials,
56+
and storage options must be supplied again when reopening.
4357

4458
Use `shared_cache=True` when separate processes need to share one `cache_dir`.
4559
They should use the same source and conversion options.

‎doc/reference/remotectable.rst‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -78,13 +78,13 @@ Parquet cache details
7878
---------------------
7979

8080
Parquet ``cache_dir`` stores each accessed physical field and row group as a
81-
native CTable directory beneath ``<source>.parquet--<hash>/<generation>.b2d/``.
82-
The generation manifest retains the table schema and row-group boundaries; a
83-
small index in ``cache_dir`` identifies its source revision. Warm opens use
84-
these local files without contacting the source. Cached sources are assumed
85-
immutable until ``refresh()`` is called. The generation path itself can be
86-
opened with :func:`blosc2.open` to recover the complete logical table and fetch
87-
uncached groups on demand.
81+
native CTable directory under the RemoteStore generation's ``parquet-groups``
82+
directory. The common manifest retains the footer, schema, conversion options,
83+
source marker, and row-group boundaries. Warm opens reuse this discovery and
84+
payload without contacting the source. Cached sources are assumed immutable
85+
until ``refresh()`` is called. New portable ``.b2z`` archives use the common
86+
RemoteStore manifest. Local Parquet sources use the same cache layout and may
87+
export ``.b2z`` references tied to the local source path.
8888

8989
HTTP Parquet sources need byte-range support. A persistent cache also needs a
9090
source size and version marker, such as an ETag, modification time, or

‎doc/reference/remotestore.rst‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
RemoteStore
44
===========
55

6-
``RemoteStore`` discovers a read-only B2Z, Zarr or HDF5 hierarchy and returns
6+
``RemoteStore`` discovers a read-only B2Z, Zarr, HDF5 or Parquet source and returns
77
:ref:`RemoteArray` and :ref:`RemoteCTable` leaves. Groups and leaves share one
88
source session: a B2Z archive, a native HDF5 index, or a Zarr store. Zarr listing
99
remains lazy.
@@ -16,6 +16,11 @@ DISK accepts ``max_cache_bytes=None`` for unbounded retention.
1616
Sources must be immutable. Generic ``blosc2.open(..., lazy=True, dataset=...)``
1717
continues to open a single array.
1818

19+
A Parquet file has one CTable at its root. Pass ``allow_table_root=True`` to
20+
retain a store handle, then use ``store[""]`` to access the table. Parquet has no
21+
child selectors. Converted physical-column/row-group caches share the store's
22+
policy, allowance, traffic counter, persistence, and refresh generation.
23+
1924
For HDF5, ``hdf5_index=`` accepts a native index dictionary, local JSON path, or
2025
remote fsspec URL. An explicit index skips hierarchy discovery and must match the
2126
source URL and selected scope. See :doc:`../guides/remote_arrays` for the

‎src/blosc2/b2objects.py‎

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -128,10 +128,6 @@ def decode_b2object_payload(payload: dict[str, Any], *, carrier_path=None, carri
128128
if carrier is None:
129129
raise ValueError("A persisted RemoteArray requires its B2ND carrier")
130130
return blosc2.RemoteArray._from_payload(payload, carrier)
131-
if kind == "remote_parquet":
132-
from blosc2.remote_parquet import RemoteParquetCTable
133-
134-
return RemoteParquetCTable._from_payload(payload)
135131
if kind == "lazyexpr":
136132
return decode_structured_lazyexpr(payload, carrier_path=carrier_path)
137133
if kind == "lazyudf":

‎src/blosc2/core.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -769,6 +769,8 @@ def parse_container_url(
769769
return urlpath, dataset, "hdf5"
770770
if any(part.endswith(".zarr") for part in parts):
771771
return urlpath, dataset, "zarr"
772+
if parsed.path.lower().endswith(".parquet"):
773+
return urlpath, dataset, "parquet"
772774

773775
return urlpath, dataset, None
774776

‎src/blosc2/remote_ctable.py‎

Lines changed: 83 additions & 84 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
import operator
1212
import os
1313

14+
import blosc2
1415
from blosc2.ctable import CTable
1516
from blosc2.ctable_storage import RemoteTableStorage
1617
from blosc2.remote_array import CACHE_POLICY_DEFAULT, RemoteMetadataMapping
@@ -110,40 +111,13 @@ def __new__(
110111
_filesystem_resolver=None,
111112
_batch_validator=None,
112113
):
113-
if source_format == "parquet" or (
114+
parquet = source_format == "parquet" or (
114115
isinstance(urlpath, (str, os.PathLike))
115116
and os.fspath(urlpath).split("?", 1)[0].lower().endswith(".parquet")
116-
):
117-
if source_format not in (None, "parquet"):
118-
raise ValueError("source_format conflicts with the .parquet suffix")
119-
from blosc2.remote_parquet import RemoteParquetCTable
120-
121-
return RemoteParquetCTable(
122-
os.fspath(urlpath),
123-
storage_options=storage_options,
124-
parquet_options=parquet_options,
125-
columns=columns,
126-
max_rows=max_rows,
127-
string_max_length=string_max_length,
128-
null_storage=null_storage,
129-
auto_null_sentinels=auto_null_sentinels,
130-
separate_nested_cols=separate_nested_cols,
131-
list_serializer=list_serializer,
132-
blosc2_batch_size=blosc2_batch_size,
133-
blosc2_items_per_block=blosc2_items_per_block,
134-
batch_size=batch_size,
135-
cparams=cparams,
136-
dparams=dparams,
137-
validate=validate,
138-
max_cache_bytes=max_cache_bytes,
139-
cache_policy=cache_policy,
140-
cache_dir=cache_dir,
141-
shared_cache=shared_cache,
142-
max_concurrency=max_concurrency,
143-
metadata_buffer_bytes=metadata_buffer_bytes,
144-
row_buffer_bytes=row_buffer_bytes,
145-
)
146-
if (
117+
)
118+
if parquet and source_format not in (None, "parquet"):
119+
raise ValueError("source_format conflicts with the .parquet suffix")
120+
if not parquet and (
147121
source_format is not None
148122
or any(
149123
value is not None
@@ -180,20 +154,61 @@ def __new__(
180154

181155
from blosc2.remote_store import RemoteStore
182156

183-
store = RemoteStore(
184-
urlpath,
185-
dataset=dataset,
186-
path=path,
187-
storage_options=storage_options,
188-
cache_policy=cache_policy,
189-
max_cache_bytes=max_cache_bytes,
190-
cache_dir=cache_dir,
191-
hdf5_index=hdf5_index,
192-
_allow_array_root=True,
193-
_filesystem=_filesystem,
194-
_filesystem_resolver=_filesystem_resolver,
195-
_batch_validator=_batch_validator,
196-
)
157+
conversion = None
158+
if parquet:
159+
conversion = {
160+
"parquet_options": parquet_options,
161+
"columns": columns,
162+
"max_rows": max_rows,
163+
"string_max_length": string_max_length,
164+
"null_storage": null_storage,
165+
"auto_null_sentinels": auto_null_sentinels,
166+
"separate_nested_cols": separate_nested_cols,
167+
"list_serializer": list_serializer,
168+
"blosc2_batch_size": blosc2_batch_size,
169+
"blosc2_items_per_block": blosc2_items_per_block,
170+
"batch_size": batch_size,
171+
"cparams": cparams,
172+
"dparams": dparams,
173+
"validate": validate,
174+
}
175+
176+
if shared_cache:
177+
if cache_dir is None or (
178+
cache_policy is not CACHE_POLICY_DEFAULT and cache_policy is not blosc2.CachePolicy.DISK
179+
):
180+
raise ValueError("shared_cache=True requires a disk cache")
181+
store = RemoteStore.with_sparse_cache(
182+
urlpath,
183+
cache_dir,
184+
dataset=dataset,
185+
path=path,
186+
storage_options=storage_options,
187+
max_cache_bytes=max_cache_bytes,
188+
_filesystem=_filesystem,
189+
_filesystem_resolver=_filesystem_resolver,
190+
_batch_validator=_batch_validator,
191+
_source_format="parquet" if parquet else None,
192+
_parquet_conversion=conversion,
193+
)
194+
else:
195+
store = RemoteStore(
196+
urlpath,
197+
dataset=dataset,
198+
path=path,
199+
storage_options=storage_options,
200+
cache_policy=cache_policy,
201+
max_cache_bytes=max_cache_bytes,
202+
cache_dir=cache_dir,
203+
hdf5_index=hdf5_index,
204+
_allow_array_root=True,
205+
_filesystem=_filesystem,
206+
_filesystem_resolver=_filesystem_resolver,
207+
_batch_validator=_batch_validator,
208+
_source_format="parquet" if parquet else None,
209+
_parquet_conversion=conversion,
210+
_allow_local_source=parquet,
211+
)
197212
try:
198213
_, full = store._resolve("")
199214
kind, diagnostic = store._owner.nodes[full]
@@ -211,12 +226,16 @@ def __init__(self, *args, **kwargs):
211226

212227
@classmethod
213228
def open_reference(cls, path, *, storage_options=None, parquet_options=None):
214-
"""Reopen a saved Parquet table with replacement runtime options."""
215-
from blosc2.remote_parquet import RemoteParquetCTable
229+
"""Reopen a saved Parquet RemoteStore archive."""
230+
if parquet_options is not None:
231+
raise TypeError("Parquet reader options are frozen in the saved reference")
232+
from blosc2.remote_store import RemoteStore
216233

217-
return RemoteParquetCTable.open_reference(
218-
path, storage_options=storage_options, parquet_options=parquet_options
219-
)
234+
table = RemoteStore._open_artifact(path, storage_options=storage_options)
235+
if not isinstance(table, cls):
236+
table.close()
237+
raise ValueError("Reference does not contain a remote CTable")
238+
return table
220239

221240
@classmethod
222241
def with_sparse_cache(
@@ -239,6 +258,7 @@ def with_sparse_cache(
239258
_source_validator=None,
240259
_manifest_validator=None,
241260
_max_nodes=None,
261+
source_format=None,
242262
):
243263
"""Attach a remote CTable to a sparse disk cache shared across processes.
244264
@@ -247,37 +267,6 @@ def with_sparse_cache(
247267
``max_cache_bytes=None`` for unlimited retention. For ordinary shared
248268
caching, prefer ``blosc2.open(url, cache_dir=..., shared_cache=True)``.
249269
"""
250-
if isinstance(urlpath, (str, os.PathLike)) and os.fspath(urlpath).split("?", 1)[0].lower().endswith(
251-
".parquet"
252-
):
253-
if any(
254-
value is not None
255-
for value in (
256-
dataset,
257-
path,
258-
manifest,
259-
carrier,
260-
_filesystem,
261-
_filesystem_resolver,
262-
_batch_validator,
263-
_source_validator,
264-
_manifest_validator,
265-
_max_nodes,
266-
)
267-
):
268-
raise ValueError("Parquet sparse caching accepts a single file and cache directory")
269-
from blosc2.remote_parquet import RemoteParquetCTable
270-
271-
return RemoteParquetCTable(
272-
os.fspath(urlpath),
273-
cache_dir=runtime_cache_path,
274-
shared_cache=True,
275-
max_cache_bytes=max_cache_bytes,
276-
storage_options=storage_options,
277-
max_concurrency=max_concurrency,
278-
metadata_buffer_bytes=metadata_buffer_bytes,
279-
row_buffer_bytes=row_buffer_bytes,
280-
)
281270
settings = {
282271
name: _positive_integer(name, value)
283272
for name, value in {
@@ -304,6 +293,7 @@ def with_sparse_cache(
304293
_source_validator=_source_validator,
305294
_manifest_validator=_manifest_validator,
306295
_max_nodes=_max_nodes,
296+
_source_format=source_format,
307297
)
308298
try:
309299
_, full = store._resolve("")
@@ -319,7 +309,14 @@ def with_sparse_cache(
319309
@classmethod
320310
def _from_owner(cls, owner, full_path, **settings):
321311
settings = {name: _positive_integer(name, value) for name, value in settings.items()}
322-
storage = RemoteTableStorage(owner, full_path, **settings)
312+
if owner.format == "parquet":
313+
from blosc2.remote_parquet import ParquetTableStorage
314+
315+
storage = ParquetTableStorage(
316+
owner, owner.parquet_schema, owner.parquet_physical, owner.parquet_length, **settings
317+
)
318+
else:
319+
storage = RemoteTableStorage(owner, full_path, **settings)
323320
try:
324321
return cls._open_from_storage(storage)
325322
except BaseException:
@@ -408,10 +405,12 @@ def attrs(self):
408405
@property
409406
def source(self):
410407
storage = self._remote_storage()
408+
from blosc2.remote_store import public_source_url
409+
411410
return {
412411
"kind": storage._owner.format,
413412
"version": 1,
414-
"urlpath": storage._owner.urlpath,
413+
"urlpath": public_source_url(storage._owner.urlpath),
415414
"dataset": storage._root_key,
416415
"assume_immutable": True,
417416
}

0 commit comments

Comments
 (0)