roboto.experimental.topics.topic#
Module Contents#
- class roboto.experimental.topics.topic.DatasetContext(/, **data)#
Bases:
pydantic.BaseModelThe dataset a Topic is scoped to: limits topic operations to the topic’s data in the dataset’s files.
An omitted
start_timeorend_timedefaults 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)
- dataset_id: str#
- model_config#
Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
- class roboto.experimental.topics.topic.DeviceContext(/, **data)#
Bases:
pydantic.BaseModelThe 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_timeorend_timedefaults 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)
- device_id: str#
- model_config#
Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
- roboto.experimental.topics.topic.FieldAddressLike#
A field-subtree address, as a
FieldAddressor explicit path components (("pose", "position")for a nested field,("angular_velocity",)for a top-level one).Each component is one
path_in_schemaelement; there is no string delimiter, so a component may itself contain a.. A bare string is rejected even though it is structurally aSequence[str]— splitting it on.would guess at component boundaries, and iterating it would address one field per character; pass the components explicitly instead.
- class roboto.experimental.topics.topic.FileContext(/, **data)#
Bases:
pydantic.BaseModelThe file a Topic is scoped to: limits topic operations to the topic’s data in that file.
An omitted
start_timeorend_timedefaults 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)
- file_id: str#
- model_config#
Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
- class roboto.experimental.topics.topic.SessionContext(/, **data)#
Bases:
pydantic.BaseModelThe 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)
- 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.
- model_config#
Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
- session_id: str#
- 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.
- roboto.experimental.topics.topic.TOPIC_DATA_CACHE_SUBDIR = 'topic-data'#
Subdirectory of the client’s cache directory where fetched topic data files are cached.
- class roboto.experimental.topics.topic.Topic(record, roboto_client=None, context=None)#
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
Topicreads org-wide and its data-returning methods require an explicit time window.A
Topiccarrying aTopicContextinstead scopes topic operations likeget_data*to the topic’s data in that context, and never reads the topic’s data elsewhere.A
SessionContext, as on a Topic yielded bylist_topics()orget_topic(), scopes to that Session’s files and defaults the window to the span of time the Session covers.A
FileContext,DatasetContextorDeviceContextscopes to that file, to that dataset’s files, or to that device’s files, and defaults the window to the start and end of the topic’s data there.
A Topic’s context is fixed when it is built.
- Parameters:
roboto_client (Optional[roboto.http.RobotoClient])
context (Optional[TopicContext])
- clear_unix_offset()#
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
TimelineExtentRecordper 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().RobotoInvalidRequestException – 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.
RobotoConflictException – A concurrent writer added files to the Session while the offset was being cleared; retry the call.
RobotoNotFoundException – The Session does not exist, or holds no data for this Topic.
RobotoUnauthorizedException – The caller cannot view the Session, or cannot edit some file behind the data being written.
- Return type:
Examples
>>> from roboto.experimental.sessions import Session >>> session = Session.from_id("se_abc123") >>> topic = session.get_topic("/camera/image_raw") >>> topic.clear_unix_offset()
- classmethod from_id(topic_id, roboto_client=None, context=None)#
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.Nonereads org-wide.
- Returns:
The loaded topic.
- Raises:
RobotoNotFoundException – No topic with this id exists.
RobotoUnauthorizedException – The caller lacks topic view access in the org that owns the topic.
- Return type:
Examples
>>> 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())
- classmethod from_record(record, roboto_client=None, context=None)#
Wrap an already-loaded topic identity record.
- Parameters:
record (roboto.domain.topics.TopicIdentityRecord) – 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.- Return type:
Examples
>>> from roboto.experimental.topics import Topic >>> topic = Topic.from_record(record) >>> topic.topic_id 'ti_abc123'
- 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)#
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.timestampis the row’s absolute Unix-epoch nanosecond timestamp (anint);recordis adictof 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:
start_time (Optional[roboto.time.Time]) – See
get_data_as_record_batches().end_time (Optional[roboto.time.Time]) – See
get_data_as_record_batches().fields_include (Optional[collections.abc.Iterable[FieldAddressLike]]) – See
get_data_as_record_batches().fields_exclude (Optional[collections.abc.Iterable[FieldAddressLike]]) – See
get_data_as_record_batches().prefer (Optional[roboto.experimental.topics.operations.RepresentationPreference]) – See
get_data_as_record_batches().schema_id (Optional[str]) – See
get_data_as_record_batches().schema_checksum (Optional[str]) – See
get_data_as_record_batches().timeline_source_id (Optional[str]) – See
get_data_as_record_batches().timeline_source_name (Optional[str]) – See
get_data_as_record_batches().cache_policy (roboto.storage.CachePolicy) – See
get_data_as_record_batches().cache_dir (Union[str, pathlib.Path, None]) – See
get_data_as_record_batches().
- Yields:
(timestamp, record)tuples for the in-window rows, filtered and projected per the arguments.- Raises:
- Return type:
collections.abc.Generator[tuple[roboto.domain.topics.Timestamp, dict[str, Any]], None, None]
Examples
>>> 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)
- 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)#
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-awareDatetimeIndex.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
flattenis set, which expands every struct level into dot-delimited leaf columns (e.g.pose.position.x). List-typed fields are unaffected byflatten.Read parameters and error behavior are as documented on
get_data_as_record_batches().Requires the
roboto[analytics]extra.- Parameters:
start_time (Optional[roboto.time.Time]) – See
get_data_as_record_batches().end_time (Optional[roboto.time.Time]) – See
get_data_as_record_batches().fields_include (Optional[collections.abc.Iterable[FieldAddressLike]]) – See
get_data_as_record_batches().fields_exclude (Optional[collections.abc.Iterable[FieldAddressLike]]) – See
get_data_as_record_batches().prefer (Optional[roboto.experimental.topics.operations.RepresentationPreference]) – See
get_data_as_record_batches().schema_id (Optional[str]) – See
get_data_as_record_batches().schema_checksum (Optional[str]) – See
get_data_as_record_batches().timeline_source_id (Optional[str]) – See
get_data_as_record_batches().timeline_source_name (Optional[str]) – See
get_data_as_record_batches().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_policy (roboto.storage.CachePolicy) – See
get_data_as_record_batches().cache_dir (Union[str, pathlib.Path, None]) – See
get_data_as_record_batches().
- Returns:
DataFrame of the in-window rows indexed by a timezone-aware
DatetimeIndex.- Raises:
- Return type:
pandas.DataFrame
Examples
>>> from roboto.experimental.topics import Topic >>> topic = Topic.from_id("ti_abc123") >>> df = topic.get_data_as_df(start_time=t0, end_time=t1)
- 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)#
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
int64column of Unix-epoch nanosecond timestamps; locate that column withtimestamp_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().Nonedefaults 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 aSessionContext(such as on a topic obtained fromlist_topics()orget_topic()); otherwise required (aValueErroris 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.
Noneprojects every field.fields_exclude (Optional[collections.abc.Iterable[FieldAddressLike]]) – Field subtrees to drop from the projection.
Nonedrops none.prefer (Optional[roboto.experimental.topics.operations.RepresentationPreference]) – Preferred representation per field subtree, selecting which stored variant of a field to read.
Noneapplies 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.
Noneuses each schema’s default source.timeline_source_name (Optional[str]) – Timeline source by name. Mutually exclusive with
timeline_source_id.cache_policy (roboto.storage.CachePolicy) – 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-datasubdirectory ofROBOTO_CACHE_DIR, or the platform-conventional per-user cache directory when that is unset.
- Yields:
pyarrow.RecordBatchinstances holding the in-window rows, filtered and projected per the arguments.- Raises:
RobotoNotFoundException – This topic’s file or dataset context names a file or dataset that does not exist.
RobotoInvalidRequestException – The window spans multiple schemas and none was chosen with
schema_idorschema_checksum, a named schema or timeline source does not match the window’s data, no stored representation satisfies a representation preference, orfields_includeandfields_excludetogether select no field. The error carries an actionable message.RobotoReadPlanExecutionException –
The topic’s data cannot be read as the service’s read plan describes it; a
RobotoInternalExceptionwhosekindsays 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_pathruns 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_numbernames 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.
RobotoInternalException – 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.RobotoUnauthorizedException – 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]
Examples
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())
- property name: str#
Human-readable topic name (e.g.
"/camera/image_raw"). Unique within an organization.- Return type:
str
- property org_id: str#
Identifier of the organization that owns this topic.
- Return type:
str
- property record: roboto.domain.topics.TopicIdentityRecord#
The underlying topic identity record.
- Return type:
- set_unix_offset(anchor)#
Anchor this Topic’s data in its Session to wall-clock time.
anchoris 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, useset_timeline_offset(); to anchor exactly the slice a declaration names, supply that declaration’sanchor_nsat ingest.Two consequences:
The data written is shared, not copied. Any other Session holding the same data reads the same anchor.
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:
anchor (roboto.time.Time) – Wall-clock instant of stored time 0: an
intof nanoseconds since the Unix epoch, or any otherTime, read asto_epoch_nanoseconds()reads it (adatetimeor ISO 8601 string is that instant; afloat,Decimal, or numeric string is seconds since the epoch). Must fall after the Unix epoch: zero is not an anchor (useclear_unix_offset()to return this Topic’s data in the Session to an offset of 0), and earlier instants are rejected.- Returns:
One
TimelineExtentRecordper timeline extent whose offset changed. An extent already anchored atanchoris 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()orlist_topics().TypeError –
anchoris not one of theTimetypes.ValueError –
anchoris 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 aspydantic.ValidationError, a subclass ofValueError.RobotoInvalidRequestException – 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.
RobotoConflictException – A concurrent writer added files to the Session while the anchor was being applied; retry the call.
RobotoNotFoundException – The Session does not exist, or holds no data for this Topic.
RobotoUnauthorizedException – The caller cannot view the Session, or cannot edit some file behind the data being written.
- Return type:
Examples
>>> 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) ... )
- property topic_id: str#
Durable identifier of this topic (
ti_*).- Return type:
str