Skip to content
Roboto
Esc
↑↓navigate↵open⌘Jpreview
On this page

roboto.formats.parquet

Fetching and decoding topic data stored in Parquet files.

Covers cache-policy-driven opening of remote Parquet files (local cache reuse, atomic download, or HTTP streaming), row-group time filtering, column projection from schema field paths, and timestamp extraction.

Submodules

Package Contents

ParquetParser

class roboto.formats.parquet.ParquetParser(source, min_required_row_group_size=100000, small_row_group_count_threshold=32)#View Source

Parameters

source pathlib.Path
min_required_row_group_size int
small_row_group_count_threshold int

Properties

ParquetParser.column_count

column_count int #
Return type: int

ParquetParser.extract_timestamp_info()

extract_timestamp_info(timestamp_column_name=None, timestamp_unit=None)#View Source

Parameters

timestamp_column_name Optional[str]
timestamp_unit Optional[Union[str, roboto.time.TimeUnit]]

Properties

ParquetParser.fields

fields Generator[pyarrow.Field, None, None] #
Return type: Generator[pyarrow.Field, None, None]

ParquetParser.find_timestamp_field_by_type()

find_timestamp_field_by_type()#View Source

Return type

pyarrow.Field

ParquetParser.get_data_for_column()

get_data_for_column(column_name)#View Source

Parameters

column_name str

Return type

pyarrow.Table

ParquetParser.get_timestamp_field_by_name()

get_timestamp_field_by_name(column_name)#View Source

Parameters

column_name str

Return type

pyarrow.Field

ParquetParser.is_parquet_file()

static is_parquet_file(path)#View Source

Parameters

path pathlib.Path

Return type

bool

ParquetParser.requires_rewrite()

requires_rewrite(timestamp)#View Source

Return type

bool

ParquetParser.rewrite()

rewrite(outfile, timestamp, target_row_group_size_bytes=100 * 1000 * 1000)#View Source

Parameters

outfile pathlib.Path
target_row_group_size_bytes int

Return type

None

Properties

ParquetParser.row_count

row_count int #
Return type: int

ParquetParser.row_group_count

row_group_count int #
Return type: int

ParquetParser.row_group_size

row_group_size int #
Return type: int

Timestamp

class roboto.formats.parquet.Timestamp#View Source

The timestamp field of a Parquet file, the field holding each row’s timestamp. Serves as both a descriptor of that field and as a utility for projecting its values to other time units.

Attributes

Timestamp.column_index

column_index int #

The timestamp’s position among the file’s leaf columns, which is the index of its column chunk in every row group.

Timestamp.field

field pyarrow.Field #

Timestamp.path

path tuple[str, ...] #

Path components locating the timestamp field in the file’s schema, root to leaf, through structs only.

Timestamp.to_epoch_nanoseconds()

to_epoch_nanoseconds(timestamp)#View Source

Parameters

Return type

int

Timestamp.unit()

unit()#View Source

Attributes

Timestamp.unit_hint

unit_hint str | None #

Unit the stored values are recorded in, used when the Arrow type does not carry one.

Taken from the Unit metadata of the timestamp’s MessagePathRecord, or from the read plan’s unit. None when the unit is unknown.

compute_time_filter_mask()

roboto.formats.parquet.compute_time_filter_mask(timestamps, start_time=None, end_time=None)#View Source

Compute a boolean mask indicating which rows fall within the specified time range. Returns None if no time filtering is needed (both start_time and end_time are None).

Parameters

timestamps pyarrow.Array
start_time Optional[int]
end_time Optional[int]

Return type

Optional[pyarrow.BooleanArray]

extract_timestamp_field()

roboto.formats.parquet.extract_timestamp_field(schema, timestamp_field, unit_hint)#View Source

Find a Parquet file’s timestamp field and describe it as a Timestamp.

The field is found by walking timestamp_field.path_in_schema one component at a time, first among the schema’s top-level fields and then through struct children, so a nested field and a top-level column whose name contains dots are never confused. schema must be the file’s ParquetFile.schema_arrow: the timestamp’s column index is counted over the schema’s leaves, which match the file’s leaf columns one for one.

unit_hint is the unit of the stored values, used when the field’s Arrow type does not carry one (an integer, floating-point or decimal field). Callers take it from the Unit metadata of the timestamp’s MessagePathRecord, or from unit.

Parameters

schema pyarrow.Schema
unit_hint Optional[str]

Raises

KeyError

A path component is not in the schema, or names a child of a field that is not a struct.

extract_timestamps()

roboto.formats.parquet.extract_timestamps(table, timestamp)#View Source

Extract timestamps in nanoseconds since Unix epoch from the table’s timestamp field.

The field is found by walking timestamp.path through the table’s struct columns. A row whose enclosing struct is null has a null timestamp.

Parameters

table pyarrow.Table

Return type

pyarrow.Int64Array

generate_message_path_requests()

roboto.formats.parquet.generate_message_path_requests(parser, timestamp, max_depth=10)#View Source

Generate AddMessagePathRequest objects for all fields in a Parquet schema.

Traverses the schema recursively to generate message paths for nested types (structs, lists) in addition to top-level fields.

Parameters

ParquetParser instance containing the schema and data.

Timestamp information for the topic.

max_depth int

Maximum recursion depth for nested types (default: 10).

Yields

AddMessagePathRequest objects for each field and nested field in the schema.

Usage

For a schema with a struct column `position: struct<x: float, y: float>`: - Yields position (Object) - Yields position.x (Number) - Yields position.y (Number)

For a schema with `values: list<float64>`: - Yields values (NumberArray)

For a schema with `points: list<struct<x: float, y: float>>`: - Yields points (Array) - Yields points.x (Number) - Yields points.y (Number)

narrow_list_nested_fields()

roboto.formats.parquet.narrow_list_nested_fields(table, schema, fields)#View Source

Prune list-of-struct columns to the projected leaves inside each element.

PyArrow’s prefix-based nested column selection cannot reach through list wrapper nodes, so resolve_columns() reads a list-nested leaf’s whole list ancestor column — every element keeps all of its struct fields. This Arrow-native post-read pass narrows each such element down to the requested leaves, leaving every other read path byte-identical.

A top-level root is narrowed iff at least one of its projected paths has a list ancestor; otherwise the table is returned unchanged (pure struct, scalar, and scalar-list reads never enter the rebuild). Per root, a trie is built from its paths with the root component stripped so non-list-nested siblings the projection also keeps are preserved. Every struct keeps its fields in the order the file stores them.

No field of fields may lie inside another: the outer one would be narrowed to the inner one rather than kept whole.

Parameters

table pyarrow.Table
schema pyarrow.Schema
fields collections.abc.Iterable[roboto.formats.fields.FieldSelection]

Return type

pyarrow.Table

open_parquet_file()

roboto.formats.parquet.open_parquet_file(url_provider, cache_outfile, policy, estimated_column_count, size_bytes=None)#View Source

Open a remote Parquet file under a cache policy, from the cheapest available source.

Dispatches on choose_fetch_mode(): an already-cached copy is reused, a download is performed (concurrency-safe, atomic) when the policy calls for one, and otherwise the file is streamed over HTTP range requests without touching disk.

Parameters

url_provider Callable[[], str]

Resolves the file’s signed download URL. Called at most once, and only when the chosen mode actually needs the URL.

cache_outfile Optional[pathlib.Path]

The file’s stable local cache path, or None when no cache location is configured (forces streaming).

The caller’s cache policy.

estimated_column_count int

How many columns the read is expected to project; informs the ADAPTIVE download-vs-stream choice.

size_bytes Optional[int]

The backing object’s size in bytes when the server reports it; None when unknown. It gates cache reuse: a present cached file whose size does not match is treated as stale and re-fetched. On the STREAM path it lets a known-large file skip the whole-file head probe; on the DOWNLOAD path it verifies the downloaded file is complete before it is promoted to the cache.

Returns

pyarrow.parquet.ParquetFile

An open pyarrow.parquet.ParquetFile.

parquet_file_from_url()

roboto.formats.parquet.parquet_file_from_url(signed_url, size_bytes=None)#View Source

Open a Parquet file over HTTP via a signed URL (no local download).

A single ranged GET over fsspec’s shared HTTP session probes the first _STREAM_WHOLE_FILE_PROBE_BYTES of the file. A file smaller than the probe arrives whole in that one request and is read from an in-memory buffer; a larger file (or a failed probe) falls back to HTTP range-request streaming through pyarrow’s filesystem layer.

When size_bytes is known and at least _STREAM_WHOLE_FILE_PROBE_BYTES, the file is known-large up front: the whole-file probe could never win (a BufferReader read is taken only for sub-threshold files), so it is skipped and the file is range-streamed directly, avoiding a wasted 16 MiB GET. size_bytes of None (an older server omits the size) preserves the probe-then-decide behavior.

Parameters

signed_url str

The file’s signed download URL.

size_bytes Optional[int]

The backing object’s size in bytes when the server reports it; None when unknown.

Raises

ValueError

The probe succeeds but the object is empty (0 bytes), which is not a readable Parquet file.

Return type

pyarrow.parquet.ParquetFile

resolve_columns()

roboto.formats.parquet.resolve_columns(schema, fields)#View Source

Build a deduplicated list of column names safe for read_row_group(columns=...).

Children of list-type columns are replaced by their list ancestor’s column name because PyArrow’s prefix-based nested column selection does not work through list wrapper nodes in the physical Parquet schema. Selecting the parent list column already returns its full nested structure.

This is important because the projected fields contain only leaf paths. For a column like points: list<struct<x, y>>, only points.x and points.y are selected — the parent points field is absent. This function derives the correct parent column name from the child’s path_in_schema.

Children of struct-type columns are preserved because PyArrow can resolve them via dot-separated prefix matching (e.g. "position.x" selects the x child of the position struct).

Each name joins path components with dots, the form in which read_row_group takes columns. Two fields whose paths join to the same name are both selected by it: "header.stamp" selects both a stamp field nested in a header struct and a top-level column of that name. select_fields() removes the fields that were not projected.

Parameters

schema pyarrow.Schema
fields collections.abc.Iterable[roboto.formats.fields.FieldSelection]

Return type

list[str]

select_fields()

roboto.formats.parquet.select_fields(table, fields)#View Source

Keep only the fields fields names, found by walking their path components, in the table’s order.

Reading columns by their dot-joined names (resolve_columns()) can bring in more than was projected, such as a top-level column named header.stamp read along with a header struct’s stamp field, or a timestamp field read only to filter rows. This removes them:

  • A top-level column that no field’s path starts with is dropped.
  • A field that names a column or struct keeps it whole, even when other fields name some of its children.
  • Otherwise a struct keeps only the children that some field’s path runs through, with its own null rows.
  • Lists and maps, and everything below them, are kept as read (narrow_list_nested_fields() narrows the structs inside a list).
  • A field whose path the table lacks is skipped.

A column from which nothing is removed is returned as read. With no fields at all, the result has no columns and keeps the table’s row count.

Parameters

table pyarrow.Table
fields collections.abc.Iterable[roboto.formats.fields.FieldSelection]

Return type

pyarrow.Table

should_narrow_list_nested_fields()

roboto.formats.parquet.should_narrow_list_nested_fields(schema, fields)#View Source

Return whether narrow_list_nested_fields() would change the table.

True iff at least one projected field addresses a leaf inside a list (its path has a list ancestor). When False, every projected field resolves through structs and scalars alone, so PyArrow’s column selection already returns the narrowed shape and the post-read prune is a no-op — callers can skip it.

Cheap enough to evaluate once per file and hoist the per-row-group narrowing decision out of the decode loop.

Parameters

schema pyarrow.Schema
fields collections.abc.Iterable[roboto.formats.fields.FieldSelection]

Return type

bool

should_read_row_group()

roboto.formats.parquet.should_read_row_group(row_group_metadata, timestamp, start_time=None, end_time=None)#View Source

Determine whether a Parquet row group contains data within the requested time range. Used to short-circuit requesting column chunks from the given row group if not relevant.

Parameters

row_group_metadata pyarrow.parquet.RowGroupMetaData
start_time Optional[int]
end_time Optional[int]

Return type

bool

timestamp_statistics()

roboto.formats.parquet.timestamp_statistics(row_group_metadata, timestamp)#View Source

The statistics of the timestamp’s column chunk in a row group, or None when there are none to use.

A nested field and a top-level column whose name contains dots can share a column chunk’s path_in_schema (header.stamp names both a header struct’s stamp field and a column of that name), so the chunk is found by its position, timestamp.column_index. Returns None when the row group has no chunk at that index, when that chunk’s path_in_schema is not the timestamp’s path joined with dots (timestamp was found in a schema other than this file’s), or when the chunk has no statistics.

Parameters

row_group_metadata pyarrow.parquet.RowGroupMetaData

Return type

Optional[pyarrow.parquet.Statistics]

Was this page helpful?