Skip to content

Commit fde7e20

Browse files
committed
Filesystem: Add reader for ORC format
1 parent 7e3d5f8 commit fde7e20

6 files changed

Lines changed: 83 additions & 0 deletions

File tree

docs/changelog.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
- Filesystem: Fail ingestion when a concrete source path matches no file while
99
keeping unmatched glob selections valid. Thanks, @hampsterx.
1010
- Filesystem: Added rsync source connector. Thanks, @oferchen.
11+
- Filesystem: Added reader for ORC format
1112

1213
## 2026/07/27 v0.8.0
1314

src/dlt_filesystem/source/adapter.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
read_jsonl,
3535
read_msgpack,
3636
read_ods,
37+
read_orc,
3738
read_parquet,
3839
read_xml,
3940
read_yaml,
@@ -78,6 +79,8 @@ def readers(
7879
filesystem_resource
7980
| dlt.transformer(name="read_msgpack", max_table_nesting=0)(read_msgpack),
8081
filesystem_resource
82+
| dlt.transformer(name="read_orc", max_table_nesting=0)(read_orc),
83+
filesystem_resource
8184
| dlt.transformer(name="read_cbor", max_table_nesting=0)(read_cbor),
8285
filesystem_resource
8386
| dlt.transformer(name="read_xml", max_table_nesting=0)(read_xml),

src/dlt_filesystem/source/format/readers.py

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,24 @@ def _polars_csv_symbols() -> Dict[str, Any]:
5252
}
5353

5454

55+
def _polars_orc_symbols() -> Dict[str, Any]:
56+
"""Symbols needed to resolve `polars.read_orc`'s type hints for casting reader hints."""
57+
58+
import fsspec
59+
import pyarrow
60+
from pandas import DataFrame
61+
from pandas._typing import DtypeBackend, FilePath, ReadBuffer
62+
63+
return {
64+
"DataFrame": DataFrame,
65+
"DtypeBackend": DtypeBackend,
66+
"FilePath": FilePath,
67+
"fsspec": fsspec,
68+
"pyarrow": pyarrow,
69+
"ReadBuffer": ReadBuffer,
70+
}
71+
72+
5573
def _polars_spreadsheet_symbols() -> Dict[str, Any]:
5674
"""Symbols needed to cast reader hint values for `polars.read_excel` and `polars.read_ods`."""
5775
from typing import Sequence
@@ -400,6 +418,31 @@ def read_spreadsheet(
400418
yield dlt.mark.with_table_name(rows, sheet_name)
401419

402420

421+
def read_orc(
422+
items: Iterator[FileItemDict],
423+
**kwargs,
424+
) -> Iterator[TDataItems]:
425+
"""Reader for ORC files."""
426+
427+
import pandas as pd
428+
429+
reader = pd.read_orc
430+
431+
kwargs = cast_kwargs_to_signature(reader, kwargs, symbols=_polars_orc_symbols())
432+
433+
for file_obj in items:
434+
with file_obj.open() as f:
435+
rec = reader(f, dtype_backend="pyarrow", **kwargs).to_records(index=False)
436+
# Turn numpy recarray record into a Python dictionary.
437+
# https://gist.github.com/rlabbe/d574eeac63fd126b2fcd1dc390cc3257
438+
# https://stackoverflow.com/a/67324508
439+
if rec.dtype is None or rec.dtype.names is None:
440+
yield rec
441+
return
442+
result = {name: rec[name] for name in rec.dtype.names}
443+
yield result
444+
445+
403446
def read_jsonl(
404447
items: Iterator[FileItemDict], chunksize: int = 1000
405448
) -> Iterator[TDataItems]:
@@ -642,6 +685,10 @@ def read_bson(self) -> DltResource:
642685
def read_msgpack(self) -> DltResource:
643686
"""MessagePack reader resource."""
644687

688+
@copy_sig(read_orc)
689+
def read_orc(self) -> DltResource:
690+
"""ORC reader resource (pyarrow)."""
691+
645692
@copy_sig(read_cbor)
646693
def read_cbor(self) -> DltResource:
647694
"""CBOR reader resource."""

src/dlt_filesystem/source/format/registry.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
"csv_headless": "read_csv_headless",
77
"jsonl": "read_jsonl",
88
"ods": "read_ods",
9+
"orc": "read_orc",
910
"parquet": "read_parquet",
1011
# bson is read-only: the file:// destination's WRITE_FORMATS is a separate tuple.
1112
"bson": "read_bson",

src/dlt_filesystem/testing/writer.py

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,9 @@
1+
import typing
2+
3+
if typing.TYPE_CHECKING:
4+
import pandas as pd
5+
6+
17
def write_bson(path, docs):
28
"""Write BSON documents concatenated into a single file (on-disk mongodump form)."""
39
import bson
@@ -27,6 +33,12 @@ def write_msgpack(path, rows, **packb_kwargs):
2733
return path
2834

2935

36+
def write_orc(path, df: "pd.DataFrame"):
37+
"""Write dataframe to ORC file."""
38+
df.to_orc(path)
39+
return path
40+
41+
3042
def write_xml(path, text):
3143
"""Write raw XML ``text`` to ``path`` as UTF-8 bytes."""
3244
with open(path, "wb") as f:
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
import pandas as pd
2+
3+
from dlt_filesystem.source.fsspec.local import LocalFilesystemSource
4+
from dlt_filesystem.testing.writer import write_orc
5+
6+
7+
def _read_via_source(path):
8+
"""Read a local ORC file end-to-end through the shared filesystem reader."""
9+
return list(LocalFilesystemSource().dlt_source(f"file://{path}", ""))
10+
11+
12+
# --- end-to-end reader (fsspec, no Docker) ---
13+
14+
15+
def test_reads_single_top_level_object(tmp_path):
16+
"""A single top-level ORC map loads as one record."""
17+
data = pd.DataFrame.from_records([{"id": 1, "name": "alice"}])
18+
path = write_orc(tmp_path / "one.orc", data)
19+
assert _read_via_source(path) == [{"id": 1, "name": "alice"}]

0 commit comments

Comments
 (0)