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

roboto.domain.topics.parquet

Submodules

Package Contents

ParquetParser

class roboto.domain.topics.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

ParquetTopicReader

class roboto.domain.topics.parquet.ParquetTopicReader(roboto_client, cache_dir=None)#View Source

Bases: roboto.domain.topics.topic_reader.TopicReader

Private interface for retrieving topic data stored in Parquet files.

Parameters

cache_dir Optional[pathlib.Path]

ParquetTopicReader.accepts()

static accepts(message_paths_to_representations)#View Source

Parameters

message_paths_to_representations collections.abc.Iterable[roboto.domain.topics.record.MessagePathRepresentationMapping]

Return type

bool

ParquetTopicReader.get_data()

get_data(message_paths_to_representations, start_time=None, end_time=None, timestamp_message_path_representation_mapping=None)#View Source

Parameters

message_paths_to_representations collections.abc.Iterable[roboto.domain.topics.record.MessagePathRepresentationMapping]
start_time Optional[int]
end_time Optional[int]
timestamp_message_path_representation_mapping Optional[roboto.domain.topics.record.MessagePathRepresentationMapping]

Return type

collections.abc.Generator[tuple[roboto.domain.topics.topic_reader.Timestamp, dict[str, Any]], None, None]

ParquetTopicReader.get_data_as_df()

get_data_as_df(message_paths_to_representations, start_time=None, end_time=None, timestamp_message_path_representation_mapping=None)#View Source

Parameters

message_paths_to_representations collections.abc.Iterable[roboto.domain.topics.record.MessagePathRepresentationMapping]
start_time Optional[int]
end_time Optional[int]
timestamp_message_path_representation_mapping Optional[roboto.domain.topics.record.MessagePathRepresentationMapping]

Return type

tuple[pandas.Series, pandas.DataFrame]

generate_message_path_requests()

roboto.domain.topics.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)

make_topic_filename_safe()

roboto.domain.topics.parquet.make_topic_filename_safe(name, replacement_char='_')#View Source

Parameters

name str
replacement_char str

Return type

str

upload_representation_file()

roboto.domain.topics.parquet.upload_representation_file(file_path, association, caller_org_id=None, roboto_client=None)#View Source

Parameters

file_path pathlib.Path
caller_org_id Optional[str]
roboto_client Optional[roboto.http.RobotoClient]

Return type

str

Was this page helpful?