Visitar URL original
Input and Output — Vortex documentation

Input and Output

Vortex arrays support reading and writing to local and remote file systems, including plain-old HTTP, S3, Google Cloud Storage, and Azure Blob Storage.

open

Lazily open a Vortex file located at the given path or URL.

open_readable

Lazily open a Vortex file through a Python object that performs the IO itself.

ReadAt

A positional byte source that Vortex reads from, for storage only reachable from Python.

ReadBytesAt

A positional byte source that returns its own buffers, for storage only reachable from Python.

SegmentCache

A segment cache that many Vortex files share, capped by total bytes.

VortexFile

Footer

The parsed footer of a Vortex file: its layout, segment map and dtype.

RepeatedScan

A prepared scan that is optimized for repeated execution.

read_url

Read a vortex struct array from a URL.

write

Write an array to a Vortex file.


vortex.open(path: str, *, store: AzureStore | CosStore | GCSStore | HfStore | HTTPStore | LocalStore | MemoryStore | S3Store | None = None, footer: Footer | None = None, without_segment_cache: bool = False, segment_cache: SegmentCache | None = None, cache_key: str | None = None) → VortexFile

Lazily open a Vortex file located at the given path or URL.

Parameters:
  • path (str) – A local path or URL to the Vortex file.

  • store – An object store created from the vortex.store package. By default the store is inferred based on the path

  • footer (vortex.file.Footer | None) – The VortexFile.footer of an earlier open of the same file. Opening then does not read the footer. Vortex cannot check that it belongs to this file, so the file must be the same one and must not have changed since.

  • without_segment_cache (bool) – If true, disable the segment cache for this file, useful when memory is constrained.

  • segment_cache (vortex.SegmentCache | None) – A cache shared with other files, used in place of this file’s own segment cache. Requires cache_key.

  • cache_key (str | None) – Identifies this file’s contents within segment_cache. Files opened with the same key share cached segments, so the key must change whenever the file does.

Examples

Open a Vortex file and perform a scan operation:

>>> import vortex as vx
>>> vxf = vx.open("data.vortex")
>>> array_iterator = vxf.scan()

See also: vortex.open_readable(), vortex.dataset.VortexDataset

vortex.open_readable(reader: ReadBytesAt | ReadAt | IO[bytes], *, footer: Footer | None = None, concurrency: int | None = None, without_segment_cache: bool = False, segment_cache: SegmentCache | None = None, cache_key: str | None = None) → VortexFile

Lazily open a Vortex file through a Python object that performs the IO itself.

Use this for storage only reachable from Python. Storage with a native object store should be opened with vortex.open() instead, which does its IO without taking the GIL.

Parameters:
  • reader (vortex.io.ReadBytesAt | vortex.io.ReadAt | binary file object) – An object implementing vortex.io.ReadBytesAt or vortex.io.ReadAt, or a binary file object with seek and readinto (or read), such as open(path, "rb"), io.BytesIO or an fsspec file. Vortex does not close it; keep it open for as long as the returned file, or anything scanned from it, is in use.

  • footer (vortex.file.Footer | None) – The VortexFile.footer of an earlier open of the same file. Opening then does no IO, which saves the footer read when a file is opened again and again. Vortex checks only that the footer fits within the size of reader, so the file must not have changed since.

  • concurrency (int | None) – The most reads to have in flight at once through a vortex.io.ReadBytesAt or vortex.io.ReadAt, 192 by default. Not accepted for a file object, whose reads are serialized because each one has to seek first.

  • without_segment_cache (bool) – If true, disable the segment cache for this file, useful when memory is constrained.

  • segment_cache (vortex.SegmentCache | None) – A cache shared with other files, used in place of this file’s own segment cache. Requires cache_key.

  • cache_key (str | None) – Identifies this file’s contents within segment_cache. Files opened with the same key share cached segments, so the key must change whenever the file does.

Examples

Open a Vortex file through an fsspec file object:

>>> import fsspec
>>> import vortex as vx
>>> with fsspec.open("memory://data.vortex", "rb") as f:
...     table = vx.open_readable(f).to_arrow().read_all()

Open the same file again without reading its footer:

>>> with fsspec.open("memory://data.vortex", "rb") as f:
...     footer = vx.open_readable(f).footer
...     vxf = vx.open_readable(f, footer=footer)
class vortex.io.ReadAt(*args, **kwargs)

A positional byte source that Vortex reads from, for storage only reachable from Python.

Pass one to vortex.open_readable(). Binary file objects (open(path, "rb"), io.BytesIO, fsspec and pyarrow files) are accepted too without implementing this protocol, but their reads are serialized because each one has to seek first. Implement ReadAt when the underlying storage supports positional reads, such as os.pread() or HTTP range requests, so that Vortex can issue many reads at once: up to 192 by default, set with the concurrency argument of vortex.open_readable().

Reads are called from Vortex worker threads, potentially many concurrently, so implementations must be safe for concurrent use. Releasing the GIL while waiting on IO lets those reads overlap.

Examples

A reader over an open file descriptor:

>>> import os
>>> class PReadFile:
...     def __init__(self, path):
...         self._fd = os.open(path, os.O_RDONLY)
...
...     def size(self) -> int:
...         return os.fstat(self._fd).st_size
...
...     def read_into(self, offset: int, buffer: memoryview) -> int:
...         return os.preadv(self._fd, [buffer], offset)
>>> vxf = vx.open_readable(PReadFile("data.vortex"))
read_into(offset: int, buffer: memoryview) → int

Read bytes starting at absolute offset into the writable buffer.

Returns the number of bytes written. A short read is retried for the remainder, and returning 0 before buffer is full is an error. buffer is valid only for the duration of the call: do not keep it, or any view derived from it, after returning.

size() → int

Total length of the source in bytes. Called once, when the file is opened.

class vortex.io.ReadBytesAt(*args, **kwargs)

A positional byte source that returns its own buffers, for storage only reachable from Python.

Pass one to vortex.open_readable(). Use it instead of ReadAt when the storage already produces a buffer, such as bytes from os.pread() or the result of an obstore range request. Vortex keeps that buffer without copying it when it is read-only, C-contiguous, the full requested length, and aligned as Vortex needs. It copies any other buffer once. An object with both read_at and read_into is read through read_at.

Reads are called from Vortex worker threads, potentially many concurrently, so implementations must be safe for concurrent use. Up to 192 reads run at once by default, set with the concurrency argument of vortex.open_readable().

Examples

A reader over an open file descriptor:

>>> import os
>>> class PReadBytes:
...     def __init__(self, path):
...         self._fd = os.open(path, os.O_RDONLY)
...
...     def size(self) -> int:
...         return os.fstat(self._fd).st_size
...
...     def read_at(self, offset: int, length: int) -> bytes:
...         return os.pread(self._fd, length, offset)
>>> vxf = vx.open_readable(PReadBytes("data.vortex"))
read_at(offset: int, length: int) → Buffer

Return up to length bytes starting at absolute offset, as any buffer-protocol object.

A shorter result is retried for the remainder, and an empty result before length bytes is an error. Vortex may keep the returned object for as long as it uses the bytes. Do not change its contents after returning it, even through a read-only view.

size() → int

Total length of the source in bytes. Called once, when the file is opened.

class vortex.SegmentCache(max_bytes)

A segment cache that many Vortex files share, capped by total bytes.

Pass it to vortex.open() or vortex.open_readable() with a cache_key. Files opened with the same key share cached segments, so a file opened again, for example once per batch, does not read them again. The key must identify the file’s contents, not only its location: a file that has changed must get a new key.

It is safe to share between threads, but it is not picklable: create one in each process.

Parameters:

max_bytes (int) – The most bytes of segments to hold. The least recently used segments are evicted first.

clear()

Remove every cached segment.

entry_count

The number of cached segments.

Return type:

int

size_bytes

The total bytes of the cached segments.

Return type:

int

final class vortex.VortexFile(file: VortexFile)
property dtype: DType

The dtype of the file.

property footer: Footer

The parsed footer, to pass to vortex.open() or vortex.open_readable() to open this file again.

property path: str

The path or URL this file was opened from.

scan(projection: Expr | list[str] | None = None, *, expr: Expr | None = None, limit: int | None = None, indices: Array | None = None, batch_size: int | None = None) → ArrayIterator

Scan the Vortex file returning a vortex.ArrayIterator.

Parameters:
  • projection (vortex.Expr | list[str] | None) – The projection expression to read, or else read all columns.

  • expr (vortex.Expr | None) – The predicate used to filter rows. The filter columns do not need to be in the projection.

  • limit (int | None) – The maximum number of rows to read after filtering. If None, read all rows.

  • indices (vortex.Array | None) – The indices of the rows to read. Must be sorted and non-null.

  • batch_size (int | None) – The number of rows to read per chunk.

Examples

Scan a file with a structured column and nulls at multiple levels and in multiple columns.

>>> import vortex as vx
>>> import vortex.expr as ve
>>> a = vx.array([
...     {'name': 'Joseph', 'age': 25},
...     {'name': None, 'age': 31},
...     {'name': 'Angela', 'age': None},
...     {'name': 'Mikhail', 'age': 57},
...     {'name': None, 'age': None},
... ])
>>> vx.io.write(a, "a.vortex")
>>> vxf = vx.open("a.vortex")
>>> vxf.scan().read_all().to_arrow_array()
<pyarrow.lib.StructArray object at ...>
-- is_valid: all not null
-- child 0 type: string_view
  [
    "Joseph",
    null,
    "Angela",
    "Mikhail",
    null
  ]
-- child 1 type: int64
  [
    25,
    31,
    null,
    57,
    null
  ]

Read just the age column:

>>> vxf.scan(['age']).read_all().to_arrow_array()
<pyarrow.lib.StructArray object at ...>
-- is_valid: all not null
-- child 0 type: int64
  [
    25,
    31,
    null,
    57,
    null
  ]

Keep rows with an age above 35. This will read O(N_KEPT) rows, when the file format allows.

>>> vxf.scan(expr=ve.column("age") > 35).read_all().to_arrow_array()
<pyarrow.lib.StructArray object at ...>
-- is_valid: all not null
-- child 0 type: string_view
  [
    "Mikhail"
  ]
-- child 1 type: int64
  [
    57
  ]
to_arrow(projection: Expr | list[str] | None = None, *, limit: int | None = None, expr: Expr | None = None, indices: Array | None = None, batch_size: int | None = None, schema: Schema | None = None) → RecordBatchReader

Scan the Vortex file as a pyarrow.RecordBatchReader.

Parameters:
  • projection (vortex.Expr | list[str] | None) – Either an expression over the columns of the file (only referenced columns will be read from the file) or an explicit list of desired columns.

  • expr (vortex.Expr | None) – The predicate used to filter rows. The filter columns need not appear in the projection.

  • indices (vortex.Array | None) – The indices of the rows to read. Must be strictly increasing and non-null.

  • batch_size (int | None) – The number of rows to read per chunk.

  • schema (pyarrow.Schema | None) – The Arrow schema to return. Use pyarrow.string() for StringArray fields. Use pyarrow.binary() for BinaryArray fields.

to_dataset() → VortexDataset

Scan the Vortex file using the pyarrow.dataset.Dataset API.

to_polars() → LazyFrame

Read the Vortex file as a pl.LazyFrame, supporting column pruning and predicate pushdown.

to_repeated_scan(projection: Expr | list[str] | None = None, *, expr: Expr | None = None, limit: int | None = None, indices: Array | None = None, batch_size: int | None = None) → RepeatedScan

Prepare a scan of the Vortex file for repeated reads, returning a vortex.RepeatedScan.

Parameters:
  • projection (vortex.Expr | list[str] | None) – The projection expression to read, or else read all columns.

  • expr (vortex.Expr | None) – The predicate used to filter rows. The filter columns do not need to be in the projection.

  • indices (vortex.Array | None) – The indices of the rows to read. Must be sorted and non-null.

  • batch_size (int | None) – The number of rows to read per chunk.

class vortex.file.Footer

The parsed footer of a Vortex file: its layout, segment map and dtype.

Pass it to vortex.open() or vortex.open_readable() to open the same file again without reading the footer. It holds no IO state, but it is not picklable.

row_count

The number of rows in the file.

Return type:

int

final class vortex.RepeatedScan(scan: RepeatedScan)

A prepared scan that is optimized for repeated execution.

execute(*, row_range: tuple[int, int] | None = None) → ArrayIterator

Execute the scan returning a vortex.ArrayIterator.

Parameters:

row_range (tuple[int, int] | None) – Tuple is interpreted as [start, stop).

Examples

Scan a file with a structured column and nulls at multiple levels and in multiple columns.

>>> import vortex as vx
>>> import vortex.expr as ve
>>> a = vx.array([
...     {'name': 'Joseph', 'age': 25},
...     {'name': None, 'age': 31},
...     {'name': 'Angela', 'age': None},
...     {'name': 'Mikhail', 'age': 57},
...     {'name': None, 'age': None},
... ])
>>> vx.io.write(a, "a.vortex")
>>> scan = vx.open("a.vortex").to_repeated_scan()
>>> scan.execute(row_range=(1, 3)).read_all().to_arrow_array()
<pyarrow.lib.StructArray object at ...>
-- is_valid: all not null
-- child 0 type: string_view
  [
    null,
    "Angela"
  ]
-- child 1 type: int64
  [
    31,
    null
  ]
scalar_at(index: int) → Scalar

Fetch a scalar from the scan returning a vortex.Scalar.

Parameters:

index (int) – The row index to fetch. Raises an IndexError if out of bounds or if the given row index was not included in the scan.

Examples

Scan a file with a structured column and nulls at multiple levels and in multiple columns.

>>> import vortex as vx
>>> import vortex.expr as ve
>>> a = vx.array([
...     {'name': 'Joseph', 'age': 25},
...     {'name': None, 'age': 31},
...     {'name': 'Angela', 'age': None},
...     {'name': 'Mikhail', 'age': 57},
...     {'name': None, 'age': None},
... ])
>>> vx.io.write(a, "a.vortex")
>>> scan = vx.open("a.vortex").to_repeated_scan()
>>> scan.scalar_at(1)
<vortex.StructScalar object at ...>
class vortex.io.VortexWriteOptions

Write Vortex files with custom configuration.

static compact()

Prioritize small size over read-throughput and read-latency.

Let’s model some stock ticker data. As you may know, the stock market always (noisly) goes up:

>>> import random
>>> sprl = vx.array([random.randint(i, i + 10) for i in range(100_000)])

If we naively wrote 8-bytes for each of these integers to a file we’d have 800,000 bytes! Let’s see how small the array buffers are when we write with the default Vortex write options (which are also used by vortex.io.write()):

>>> vx.io.VortexWriteOptions.default().write(sprl, "chonky.vortex")
>>> vx.open("chonky.vortex").scan().read_all().nbytes
213248

Wow, Vortex manages to use about two bytes per integer! So advanced. So tiny.

But can we do better?

We sure can.

>>> vx.io.VortexWriteOptions.compact().write(sprl, "tiny.vortex")
>>> vx.open("tiny.vortex").scan().read_all().nbytes
52564

Random numbers are not (usually) composed of random bytes!

static default()

Balance size, read-throughput, and read-latency.

write(iter, path, *, store=None)

Write an array or iterator of arrays to a file.

Parameters:

Examples

Write a single Vortex array a to the local file a.vortex using the default settings:

>>> import vortex as vx
>>> import random
>>> a = vx.array([0, 1, 2, 3, None, 4])
>>> vx.io.VortexWriteOptions.default().write(a, "a.vortex")

Write the same array while preferring small file sizes over read-throughput and read-latency:

>>> import vortex as vx
>>> vx.io.VortexWriteOptions.compact().write(a, "a.vortex")
vortex.io.read_url(url, *, store=None, projection=None, row_filter=None, indices=None, row_range=None)

Read a vortex struct array from a URL.

Parameters:

Examples

Read an array from an HTTPS URL:

>>> import vortex as vx
>>> a = vx.io.read_url("https://example.com/dataset.vortex")

Read an array from an S3 URL:

>>> a = vx.io.read_url("s3://bucket/path/to/dataset.vortex")

Read an array from an Azure Blob File System URL:

>>> a = vx.io.read_url("abfss://my_file_system@my_account.dfs.core.windows.net/path/to/dataset.vortex")

Read an array from an Azure Blob Storage URL:

>>> a = vx.io.read_url("https://my_account.blob.core.windows.net/my_container/path/to/dataset.vortex")

Read an array from a Google Storage URL:

>>> a = vx.io.read_url("gs://bucket/path/to/dataset.vortex")

Read an array from a local file URL:

>>> a = vx.io.read_url("file:///path/to/dataset.vortex")

Read from S3 with explicit credentials:

>>> from vortex import store as S
>>> store = S.S3Store(
...     bucket="my-bucket",
...     region="us-east-1",
...     access_key_id="AKIA...",
...     secret_access_key="..."
... )
>>> a = vx.io.read_url("s3://my-bucket/data.vortex", store=store)
vortex.io.write(iter, path, *, store=None)

Write an array to a Vortex file.

Parameters:

Examples

Write a single Vortex array a to the local file a.vortex.

>>> import vortex as vx
>>> a = vx.array([
...     {'x': 1},
...     {'x': 2},
...     {'x': 10},
...     {'x': 11},
...     {'x': None},
... ])
>>> vx.io.write(a, "a.vortex")

Stream a PyArrow Table directly to Vortex without loading into memory:

>>> import pyarrow as pa
>>> import vortex as vx
>>> table = pa.table({'x': [1, 2, 3, 4, 5]})
>>> vx.io.write(table, "streamed.vortex")

Stream from a PyArrow RecordBatchReader:

>>> import pyarrow as pa
>>> import vortex as vx
>>> reader = pa.RecordBatchReader.from_batches(schema, batches)
>>> vx.io.write(reader, "streamed.vortex")