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.
Lazily open a Vortex file located at the given path or URL. |
|
Lazily open a Vortex file through a Python object that performs the IO itself. |
|
A positional byte source that Vortex reads from, for storage only reachable from Python. |
|
A positional byte source that returns its own buffers, for storage only reachable from Python. |
|
A segment cache that many Vortex files share, capped by total bytes. |
|
The parsed footer of a Vortex file: its layout, segment map and dtype. |
|
A prepared scan that is optimized for repeated execution. |
|
Read a vortex struct array from a URL. |
|
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) – TheVortexFile.footerof 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. Requirescache_key.cache_key (
str| None) – Identifies this file’s contents withinsegment_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 implementingvortex.io.ReadBytesAtorvortex.io.ReadAt, or a binary file object withseekandreadinto(orread), such asopen(path, "rb"),io.BytesIOor 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) – TheVortexFile.footerof 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 ofreader, so the file must not have changed since.concurrency (
int| None) – The most reads to have in flight at once through avortex.io.ReadBytesAtorvortex.io.ReadAt, 192 by default. Not accepted for a file object, whose reads are serialized because each one has toseekfirst.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. Requirescache_key.cache_key (
str| None) – Identifies this file’s contents withinsegment_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 toseekfirst. ImplementReadAtwhen the underlying storage supports positional reads, such asos.pread()or HTTP range requests, so that Vortex can issue many reads at once: up to 192 by default, set with theconcurrencyargument ofvortex.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
offsetinto the writablebuffer.Returns the number of bytes written. A short read is retried for the remainder, and returning
0beforebufferis full is an error.bufferis valid only for the duration of the call: do not keep it, or any view derived from it, after returning.
- 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 ofReadAtwhen the storage already produces a buffer, such asbytesfromos.pread()or the result of anobstorerange 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 bothread_atandread_intois read throughread_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
concurrencyargument ofvortex.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
lengthbytes starting at absoluteoffset, as any buffer-protocol object.A shorter result is retried for the remainder, and an empty result before
lengthbytes 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.
- class vortex.SegmentCache(max_bytes)¶
A segment cache that many Vortex files share, capped by total bytes.
Pass it to
vortex.open()orvortex.open_readable()with acache_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.
- final class vortex.VortexFile(file: VortexFile)¶
-
The parsed footer, to pass to
vortex.open()orvortex.open_readable()to open this file again.
- 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. Usepyarrow.string()forStringArrayfields. Usepyarrow.binary()forBinaryArrayfields.
- to_dataset() VortexDataset¶
Scan the Vortex file using the
pyarrow.dataset.DatasetAPI.
- 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.
The parsed footer of a Vortex file: its layout, segment map and dtype.
Pass it to
vortex.open()orvortex.open_readable()to open the same file again without reading the footer. It holds no IO state, but it is not picklable.The number of rows in the file.
- Return type:
- 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.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
IndexErrorif 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.
See also
- 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:
iter (vortex.Array | vortex.ArrayIterator | pyarrow.Table | pyarrow.RecordBatchReader) – The data to write. Can be a single array, an array iterator, or a PyArrow object that supports streaming. When using PyArrow objects, data is streamed directly without loading the entire dataset into memory.
path (str) – The file path.
store (vortex.store.AzureStore | vortex.store.CosStore | vortex.store.GCSStore | vortex.store.HfStore | vortex.store.HTTPStore | vortex.store.LocalStore | vortex.store.MemoryStore | vortex.store.S3Store | None) – An optional object store configuration to use for writing the output.
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")
See also
- 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:
url (str) – The URL to read from.
store (vortex.store.AzureStore | vortex.store.CosStore | vortex.store.GCSStore | vortex.store.HfStore | vortex.store.HTTPStore | vortex.store.LocalStore | vortex.store.MemoryStore | vortex.store.S3Store | None) – Pre-configured object store with credentials and settings. If provided, uses this store’s configuration. If None, checks session registry for matching URL pattern. If not found, raises VortexError.
projection (list[str | int] | None) – The columns to read identified either by their index or name.
row_filter (Expr | None) – Keep only the rows for which this expression evaluates to true.
indices (Array | None) – A list of rows to keep identified by the zero-based index within the file. NB: If row_range is specified, these indices are within the row range, not the file!
row_range (tuple[int, int] | None) – A left-inclusive, right-exclusive range of rows to read.
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:
iter (vortex.Array | vortex.ArrayIterator | pyarrow.Table | pyarrow.RecordBatchReader) – The data to write. Can be a single array, an array iterator, or a PyArrow object that supports streaming. When using PyArrow objects, data is streamed directly without loading the entire dataset into memory.
path (str) – The file path.
store (vortex.store.AzureStore | vortex.store.CosStore | vortex.store.GCSStore | vortex.store.HfStore | vortex.store.HTTPStore | vortex.store.LocalStore | vortex.store.MemoryStore | vortex.store.S3Store | None) – An optional object store configuration to use for writing the output.
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")
See also