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

roboto.experimental.topics.topic

Module Contents

DatasetContext

class roboto.experimental.topics.topic.DatasetContext(/, **data)#View Source

Bases: pydantic.BaseModel

The dataset a Topic is scoped to: limits topic operations to the topic’s data in the dataset’s files.

An omitted start_time or end_time defaults to the start or end of the topic’s data in those files, resolved by the service when the data is read.

Parameters

data Any

Attributes

DatasetContext.dataset_id

dataset_id str #

DatasetContext.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

DeviceContext

class roboto.experimental.topics.topic.DeviceContext(/, **data)#View Source

Bases: pydantic.BaseModel

The device a Topic is scoped to: limits topic operations to the topic’s data in the device’s files.

A file belongs to the device it names, or, when it names none, to the device its dataset names. An omitted start_time or end_time defaults to the start or end of the topic’s data in those files, resolved by the service when the data is read.

Parameters

data Any

Attributes

DeviceContext.device_id

device_id str #

DeviceContext.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

FieldAddressLike

roboto.experimental.topics.topic.FieldAddressLike#View Source

A field-subtree address, as a FieldAddress or explicit path components (("pose", "position") for a nested field, ("angular_velocity",) for a top-level one).

Each component is one path_in_schema element; there is no string delimiter, so a component may itself contain a .. A bare string is rejected even though it is structurally a Sequence[str] — splitting it on . would guess at component boundaries, and iterating it would address one field per character; pass the components explicitly instead.

FileContext

class roboto.experimental.topics.topic.FileContext(/, **data)#View Source

Bases: pydantic.BaseModel

The file a Topic is scoped to: limits topic operations to the topic’s data in that file.

An omitted start_time or end_time defaults to the start or end of the topic’s data in the file, resolved by the service when the data is read.

Parameters

data Any

Attributes

FileContext.file_id

file_id str #

FileContext.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

SessionContext

class roboto.experimental.topics.topic.SessionContext(/, **data)#View Source

Bases: pydantic.BaseModel

The Session a Topic is scoped to: limits topic operations to the Session’s associated files and supplies the Session’s aggregate time window as the default window for those operations.

Parameters

data Any

Attributes

SessionContext.end_time

end_time int | None = None #

Latest time covered by the Session (Unix-epoch ns); the default end_time for get_data*. None when the Session includes no files.

SessionContext.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

SessionContext.session_id

session_id str #

SessionContext.start_time

start_time int | None = None #

Earliest time covered by the Session (Unix-epoch ns); the default start_time for get_data*. None when the Session includes no files.

TOPIC_DATA_CACHE_SUBDIR

roboto.experimental.topics.topic.TOPIC_DATA_CACHE_SUBDIR = 'topic-data'#View Source

Subdirectory of the client’s cache directory where fetched topic data files are cached.

Topic

class roboto.experimental.topics.topic.Topic(record, roboto_client=None, context=None)#View Source

A logical stream of robotics data, identified durably across the files that carry it.

Within an organization, topic names are unique; contributions from different files with the same topic name share a single topic identity. By default a Topic reads org-wide and its data-returning methods require an explicit time window.

A Topic carrying a TopicContext instead scopes topic operations like get_data* to the topic’s data in that context, and never reads the topic’s data elsewhere.

A Topic’s context is fixed when it is built.

Parameters

Topic.clear_unix_offset()

clear_unix_offset()#View Source

Return this Topic’s data in its Session to an offset of 0.

Its stored timestamps are then read as nanoseconds since the Unix epoch with nothing added. This call reaches the same data as set_unix_offset(), and a file’s declared time range in the Session moves back under the same condition: only when everything that range covers moved by the same distance.

Returns

One TimelineExtentRecord per timeline extent whose offset changed. An extent already at an offset of 0 is left untouched and absent from the list, so an empty list means there was no anchor to clear.

Raises

ValueError

This Topic carries no Session context; see set_unix_offset().

Some of this Topic’s own timestamps are negative, so an offset of 0 would place that data before the Unix epoch; or a time range declared over the data would start before the epoch once moved back with it.

A concurrent writer added files to the Session while the offset was being cleared; retry the call.

The Session does not exist, or holds no data for this Topic.

The caller cannot view the Session, or cannot edit some file behind the data being written.

Usage

from roboto.experimental.sessions import Session
session = Session.from_id("se_abc123")
topic = session.get_topic("/camera/image_raw")
topic.clear_unix_offset()

Topic.from_id()

classmethod from_id(topic_id, roboto_client=None, context=None)#View Source

Load an existing topic by its id.

Parameters

topic_id str

Identifier of the topic (ti_*).

roboto_client Optional[roboto.http.RobotoClient]

Roboto client instance. Uses the default if omitted.

context Optional[TopicContext]

Optional. When provided, limits topic operations to the topic’s data in one Session, file, dataset or device, and supplies the default read window, as described on Topic. None reads org-wide.

Returns

The loaded topic.

Raises

No topic with this id exists.

The caller lacks topic view access in the org that owns the topic.

Usage

from roboto.experimental.topics import Topic
topic = Topic.from_id("ti_abc123")
topic.name
# '/camera/image_raw'

Read all of the topic’s data in one file:

from roboto.experimental.topics import FileContext
file_topic = Topic.from_id("ti_abc123", context=FileContext(file_id="fl_abc123"))
rows = list(file_topic.get_data())

Topic.from_record()

classmethod from_record(record, roboto_client=None, context=None)#View Source

Wrap an already-loaded topic identity record.

Parameters

The topic identity record to wrap.

roboto_client Optional[roboto.http.RobotoClient]

Roboto client instance. Uses the default if omitted.

context Optional[TopicContext]

Optional scope; see from_id().

Returns

A topic backed by record, with no further service calls.

Usage

from roboto.experimental.topics import Topic
topic = Topic.from_record(record)
topic.topic_id
# 'ti_abc123'

Topic.get_data()

get_data(start_time=None, end_time=None, fields_include=None, fields_exclude=None, prefer=None, schema_id=None, schema_checksum=None, timeline_source_id=None, timeline_source_name=None, cache_policy=CachePolicy.ADAPTIVE, cache_dir=None)#View Source

Yield this topic’s data within a time window, as (timestamp, record) pairs.

Convenience over get_data_as_record_batches() that unpacks each Arrow RecordBatch into one (timestamp, record) tuple per row. timestamp is the row’s absolute Unix-epoch nanosecond timestamp (an int); record is a dict of the projected fields, with struct fields as nested dicts and list fields as lists. A field the data omits for a row is absent from (or null within) that row’s dict.

Time windowing, field projection, representation selection, sort order, and error behavior are all as documented on get_data_as_record_batches().

Requires the roboto[analytics] extra.

Parameters

fields_include Optional[collections.abc.Iterable[FieldAddressLike]]
fields_exclude Optional[collections.abc.Iterable[FieldAddressLike]]
schema_id Optional[str]
schema_checksum Optional[str]
timeline_source_id Optional[str]
timeline_source_name Optional[str]
cache_dir Union[str, pathlib.Path, None]

Yields

(timestamp, record) tuples for the in-window rows, filtered and projected per the arguments.

Return type

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

Usage

from roboto.experimental.topics import Topic
topic = Topic.from_id("ti_abc123")
for timestamp, record in topic.get_data(start_time=t0, end_time=t1):
    print(timestamp, record)

Topic.get_data_as_df()

get_data_as_df(start_time=None, end_time=None, fields_include=None, fields_exclude=None, prefer=None, schema_id=None, schema_checksum=None, timeline_source_id=None, timeline_source_name=None, flatten=False, cache_policy=CachePolicy.ADAPTIVE, cache_dir=None)#View Source

Return this topic’s data within a time window as a pandas DataFrame.

Same pipeline as get_data_as_record_batches(), with the batches packed into a DataFrame whose index is a timezone-aware DatetimeIndex.

Rows return ordered by partition (each file’s data in start order), not interleaved across partitions; within a partition rows keep their stored order. Call df.sort_index() for a strict row-level time-ordered view.

A struct field is returned as a single schema-shaped column of dicts unless flatten is set, which expands every struct level into dot-delimited leaf columns (e.g. pose.position.x). List-typed fields are unaffected by flatten.

Read parameters and error behavior are as documented on get_data_as_record_batches().

Requires the roboto[analytics] extra.

Parameters

fields_include Optional[collections.abc.Iterable[FieldAddressLike]]
fields_exclude Optional[collections.abc.Iterable[FieldAddressLike]]
schema_id Optional[str]
schema_checksum Optional[str]
timeline_source_id Optional[str]
timeline_source_name Optional[str]
flatten bool

Expand struct-typed fields into dot-delimited leaf columns. When False, each struct-typed field is a single object-dtype column of dicts.

cache_dir Union[str, pathlib.Path, None]

Returns

pandas.DataFrame

DataFrame of the in-window rows indexed by a timezone-aware DatetimeIndex.

Usage

from roboto.experimental.topics import Topic
topic = Topic.from_id("ti_abc123")
df = topic.get_data_as_df(start_time=t0, end_time=t1)

Topic.get_data_as_record_batches()

get_data_as_record_batches(start_time=None, end_time=None, fields_include=None, fields_exclude=None, prefer=None, schema_id=None, schema_checksum=None, timeline_source_id=None, timeline_source_name=None, cache_policy=CachePolicy.ADAPTIVE, cache_dir=None)#View Source

Yield this topic’s data within a time window, as Arrow RecordBatches.

Each batch carries one column per top-level projected field, with nested struct and list types mirroring the topic’s schema, pruned to the projection, plus a dedicated int64 column of Unix-epoch nanosecond timestamps; locate that column with timestamp_column_index(). A field the data omits for a row surfaces as null at the deepest level that represents the omission (a whole absent subtree is a single null).

Batch sizes and boundaries carry no meaning, and a window matching no rows yields no batches, as does a context holding no data for the topic. A topic’s data can span several files (“topic partitions”); rows from different partitions are never mixed within a batch. Partitions arrive ordered by when each file’s data begins. Within a partition, rows keep their stored order, and rows from different partitions are never interleaved. So batches arrive as whole partitions in start order, not as a globally time-sorted row stream. Sort downstream if a strict row-level time order is needed.

Requires the roboto[analytics] extra.

Parameters

start_time Optional[roboto.time.Time]

Inclusive window lower bound, as nanoseconds since the Unix epoch or anything convertible via to_epoch_nanoseconds(). None defaults to the start of the topic’s data in this topic’s file, dataset or device context, or to the earliest time the Session covers in a SessionContext (such as on a topic obtained from list_topics() or get_topic()); otherwise required (a ValueError is raised when it cannot be resolved). An explicit bound narrows the read within the context and never widens it.

end_time Optional[roboto.time.Time]

Inclusive window upper bound, same forms as start_time; defaults to the end of the topic’s data, or to the latest time the Session covers, on the same terms.

fields_include Optional[collections.abc.Iterable[FieldAddressLike]]

Field subtrees to project. None projects every field.

fields_exclude Optional[collections.abc.Iterable[FieldAddressLike]]

Field subtrees to drop from the projection. None drops none.

Preferred representation per field subtree, selecting which stored variant of a field to read. None applies the default selection everywhere.

schema_id Optional[str]

Schema to read under, by id. Required only when the window spans data with more than one schema.

schema_checksum Optional[str]

Schema to read under, by checksum. Mutually exclusive with schema_id.

timeline_source_id Optional[str]

Timeline source to resolve the window with, by id. None uses each schema’s default source.

timeline_source_name Optional[str]

Timeline source by name. Mutually exclusive with timeline_source_id.

Whether fetched Parquet files are cached to local disk. MCAP data always streams.

cache_dir Union[str, pathlib.Path, None]

Directory topic data files are cached under. Defaults to a topic-data subdirectory of ROBOTO_CACHE_DIR, or the platform-conventional per-user cache directory when that is unset.

Yields

pyarrow.RecordBatch instances holding the in-window rows, filtered and projected per the arguments.

Raises

This topic’s file or dataset context names a file or dataset that does not exist.

The window spans multiple schemas and none was chosen with schema_id or schema_checksum, a named schema or timeline source does not match the window’s data, no stored representation satisfies a representation preference, or fields_include and fields_exclude together select no field. The error carries an actionable message.

The topic’s data cannot be read as the service’s read plan describes it; a RobotoInternalException whose kind says why:

  • field-not-in-file: a file backing this topic lacks a field the read takes from it, such as a projected field or the field holding each row’s timestamp. field_path runs from that field’s top-level field down to the first component the file lacks.
  • unsupported-timestamp: the read cannot take timestamps from where the data keeps them, such as a message time on a Parquet file, or a Parquet timestamp field that is not a number.
  • invalid-timestamp: a stored timestamp is NaN or infinite, is a DECIMAL holding a fraction of a nanosecond, or leaves the signed 64-bit range once shifted to absolute time. row_number names the row by its 0-based position among the topic’s rows in its file.
  • scan-task-row-mismatch: files that store different fields of the same rows hold different rows in the window.
  • field-split-inside-non-struct: files that store different fields of the same rows split a field that one of them stores as other than a struct, such as a list or a map.
  • projected-field-in-no-scan-task: no scan task of a topic partition reads the whole schema or a subtree containing a projected field.
  • inconsistent-scan-tasks-on-file: the plan reads one file two ways, with a different format, transformations or topic name.
  • partition-schema-mismatch: a topic partition’s files give the read a different schema than the first topic partition’s files. It is raised when that partition’s files open, before any of its rows, so even when it has no rows in the window.
  • plan-without-schema: the plan reads every field of its schema and has a topic partition with a scan task, but names no schema.
  • data-range-not-in-file: a topic partition’s declared slice of its file (data_range) ends past the file’s stored row count, so the slice does not match the file.

An MCAP file backing this topic cannot be decoded, such as one with no chunk or message index, or one whose messages are protobuf-encoded.

The caller lacks read access to at least one in-window file backing this topic, or cannot view the file or dataset this topic’s context names.

Return type

collections.abc.Generator[pyarrow.RecordBatch, None, None]

Usage

Print every record in a window:

from roboto.experimental.topics import Topic
topic = Topic.from_id("ti_abc123")
for batch in topic.get_data_as_record_batches(start_time=t0, end_time=t1):
    print(batch.num_rows, batch.schema.names)

Project to one field subtree, dropping one of its children:

for batch in topic.get_data_as_record_batches(
    start_time=t0,
    end_time=t1,
    fields_include=[("angular_velocity",)],
    fields_exclude=[("angular_velocity", "y")],
):
    print(batch.to_pylist())

Properties

Topic.name

name str #

Human-readable topic name (e.g. "/camera/image_raw"). Unique within an organization.

Return type: str

Topic.org_id

org_id str #

Identifier of the organization that owns this topic.

Return type: str

Topic.record

The underlying topic identity record.

Topic.set_unix_offset()

set_unix_offset(anchor)#View Source

Anchor this Topic’s data in its Session to wall-clock time.

anchor is the wall-clock instant at which stored time 0 of that data occurred. Converted to nanoseconds since the Unix epoch, it is added to the stored timestamps of every part of this Topic the Session holds, so each timestamp reads as wall-clock time. Use this call for a Session whose data all starts from one time 0 but is stored apart: a topic chunked across several files, or several Sessions packed into slices of one shared file where this Session holds one of the slices. Every part moves in one transaction, so a failure leaves every one of them at the offset it already had.

This call covers one Topic within one Session. To anchor everything in a Session, use set_unix_offset(); to anchor a whole file regardless of Session, use set_timeline_offset(); to anchor exactly the slice a declaration names, supply that declaration’s anchor_ns at ingest.

Two consequences:

  1. The data written is shared, not copied. Any other Session holding the same data reads the same anchor.
  2. A file’s declared time range in the Session moves only when everything that range covers moved by the same distance. A file whose range also covers another topic’s data, which this call leaves alone, keeps the range it has.

Parameters

Wall-clock instant of stored time 0: an int of nanoseconds since the Unix epoch, or any other Time, read as to_epoch_nanoseconds() reads it (a datetime or ISO 8601 string is that instant; a float, Decimal, or numeric string is seconds since the epoch). Must fall after the Unix epoch: zero is not an anchor (use clear_unix_offset() to return this Topic’s data in the Session to an offset of 0), and earlier instants are rejected.

Returns

One TimelineExtentRecord per timeline extent whose offset changed. An extent already anchored at anchor is left untouched and absent from the list, so an empty list means the anchor was already in place.

Raises

ValueError

This Topic carries no Session context, so there is no way to tell which of the Topic’s data is meant; reach it through get_topic() or list_topics().

TypeError

anchor is not one of the Time types.

ValueError

anchor is a boolean, a string that is neither seconds nor ISO 8601, zero, before the Unix epoch, or too large for a signed 64-bit integer of nanoseconds; rejected client-side, before any request is made. A range refusal is raised as pydantic.ValidationError, a subclass of ValueError.

The anchor would move this Topic’s data, or a time range declared over it, before the Unix epoch or past the largest storable Unix-epoch nanosecond value; anchor the data at the instant it was recorded.

A concurrent writer added files to the Session while the anchor was being applied; retry the call.

The Session does not exist, or holds no data for this Topic.

The caller cannot view the Session, or cannot edit some file behind the data being written.

Usage

from roboto.experimental.sessions import Session
session = Session.from_id("se_abc123")
topic = session.get_topic("/camera/image_raw")
topic.set_unix_offset(1_700_000_000_000_000_000)

The same anchor given as a datetime:

import datetime
topic.set_unix_offset(
    datetime.datetime(2023, 11, 14, 22, 13, 20, tzinfo=datetime.timezone.utc)
)

Properties

Topic.topic_id

topic_id str #

Durable identifier of this topic (ti_*).

Return type: str

TopicContext

roboto.experimental.topics.topic.TopicContext#View Source

The scope a Topic reads within: one Session, file, dataset or device.

Was this page helpful?