roboto.experimental.topics
Topics APIs in active refinement; see roboto.experimental for the stability contract.
Submodules
- roboto.experimental.topics.batch_transforms
- roboto.experimental.topics.decode
- roboto.experimental.topics.operations
- roboto.experimental.topics.plan_execution
- roboto.experimental.topics.read_plan
- roboto.experimental.topics.record
- roboto.experimental.topics.split_partition
- roboto.experimental.topics.topic
Package Contents
DatasetContext
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 AnyDeviceContext
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 AnyFieldAddress
Bases: pydantic.BaseModel
Addresses a schema field, and the subtree nested under it, by exactly one of two forms.
A path names the field by its path_in_schema components directly (no string delimiter, so a component may itself contain a .); a field_id names it opaquely and resolves server-side to the same path. Either form designates the field and every field nested under it.
Parameters
data AnyAttributes
FieldAddress.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
FieldAddress.path
The field’s path_in_schema components; () addresses the schema root.
FieldAddressLike
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
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 AnyPLAN_VERSION
Contract version stamped on every plan.
ReadPlan validation refuses a plan whose version it does not recognize, so a consumer on an older contract fails at parse time instead of misreading a newer plan.
ReadPlan
Bases: pydantic.BaseModel
Resolves a read of one topic over a time window into the files to fetch and how to interpret them.
Parameters
data AnyAttributes
ReadPlan.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlan.partitions
One entry per partition in the window, each its own fetch-and-interpret plan.
Ordered by when each file’s data begins. Partitions backed by the same file are contiguous in this tuple, ordered by where their data_range slice starts in that file, so they arrive in the order their data sits in it. This orders whole partitions, not rows.
ReadPlan.plan_version
Contract version of this plan. Validation refuses a version this model does not recognize.
ReadPlan.projection
The output fields a consumer projects decoded rows to.
ReadPlan.schema_
The schema the plan reads under. Serializes as schema.
None when no partition of the topic lies in the window (within the request’s restrictions), and the plan then has no partitions. A plan that names a schema may also have no partitions.
ReadPlan.window
The time window the caller asked to read, which the plan was resolved over.
A consumer selects rows by each partition’s ReadPlanPartition.window alone, which can be narrower than this one.
None only when the request omitted a bound and either its restrictions select no timestamped data for the topic, or the bound it gave lies past the far end of that data (a start after the data ends, or an end before it begins); the plan then has no partitions.
ReadPlanFieldRef
Bases: pydantic.BaseModel
A schema field named by both its id and its path components in the schema.
Parameters
data AnyAttributes
ReadPlanFieldRef.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanFieldRef.path
The field’s path components within the schema, from the root to the field.
ReadPlanObjectRef
Bases: pydantic.BaseModel
Points to the file backing a scan task. A consumer fetches the file’s bytes from it.
Parameters
data AnyAttributes
ReadPlanObjectRef.fs_node_id
Identifier of the backing file. This id is stable, so a consumer can cache on it.
ReadPlanObjectRef.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanPartition
Bases: pydantic.BaseModel
Everything needed to fetch and interpret one in-window partition’s bytes.
Parameters
data AnyAttributes
ReadPlanPartition.data_range
The slice of the file this partition owns, as (start, end), or None for the whole file.
start is the first owned position and end is one past the last, like a Python slice. Positions are in the addressing native to the file the slice was declared against: stored-row positions counted from 0, or nanoseconds of media time for video. Any slice a consumer decodes is a range of stored rows.
A slice appears when one file packs several partitions’ data (a LeRobot v3 data file, for example). Each of those partitions carries its own time_offset_ns and may cover the same instants as the partitions beside it in the file, so time alone cannot separate its rows from theirs. Restrict each scan task’s decode to this slice before filtering by time, or the read returns the neighboring partitions’ rows too.
A scan task may name a transformed representation of the file the slice was declared against (see ReadPlanScanTask.transformations) rather than that file itself. The slice still applies unchanged: a partition that owns a slice is served only by representations holding the same content at the same positions.
ReadPlanPartition.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanPartition.scan_tasks
The files to read for this partition; empty when the partition has no readable data.
A partition may have several scan tasks, each covering a subtree of the schema or the whole schema; the read takes each field from the highest-precedence scan task whose subtree contains it.
ReadPlanPartition.time_offset_ns
Offset a consumer adds to each decoded row timestamp; the same for every row in the partition.
ReadPlanPartition.timestamp
Where this partition’s row timestamps come from.
ReadPlanPartition.window
The window this partition’s rows are selected in, in absolute Unix-epoch nanoseconds.
A consumer keeps only the rows whose stored timestamp, once time_offset_ns is added, falls inside this window, inclusive on both ends. It is the only time filter a read applies.
It is the caller’s window intersected with the part of this partition’s backing file the read is scoped to, so it can be narrower than ReadPlan.window: a read scoped to a Session narrows it when the Session holds that file over only part of the file’s time span. The narrowing is stated per file, so every partition backed by one file carries the same window; ReadPlanPartition.data_range separates the partitions packed into one file.
ReadPlanProjection
Bases: pydantic.BaseModel
The output fields the plan resolves rows to.
The projection takes exactly one of two forms: either every field in the schema (all is true, and the field list is left implicit so the plan need not enumerate a large schema) or an explicit fields list.
Parameters
data AnyAttributes
ReadPlanProjection.all_fields()
Return a projection covering every field in the schema.
Return type
Attributes
ReadPlanProjection.fields
The resolved field set when the read is narrowed; None when all is true.
ReadPlanProjection.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanProjection.narrowed()
Return a projection narrowed to an explicit field set.
Parameters
fields Iterable[ReadPlanFieldRef]Return type
ReadPlanRequest
Bases: pydantic.BaseModel
The body of a read-plan request: the logical read question to resolve into a physical plan.
session_id, file_id, dataset_id and device_id are restrictions: each limits the read to the topic’s data in one session, file, dataset or device. Every restriction the request names narrows the read, and restrictions named together intersect.
Parameters
data AnyAttributes
ReadPlanRequest.dataset_id
Limits the read to the topic’s data in this dataset’s files; None adds no dataset restriction.
ReadPlanRequest.device_id
Limits the read to the topic’s data in this device’s files; None adds no device restriction.
A file belongs to the device it names, or, when it names none, to the device its dataset names.
ReadPlanRequest.end_time
Inclusive window upper bound, absolute Unix-epoch nanoseconds.
May be None on the same terms as start_time; the bound then defaults to the latest time of the topic’s data within the request’s restrictions.
ReadPlanRequest.fields_exclude
Field subtrees to drop from the projection; None drops none.
ReadPlanRequest.fields_include
Field subtrees to project; None projects every field.
ReadPlanRequest.file_id
Limits the read to the topic’s data in this file; None adds no file restriction.
ReadPlanRequest.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanRequest.prefer
Per-subtree representation preference; None applies default selection everywhere.
ReadPlanRequest.schema_checksum
Schema to use, by checksum, or None.
ReadPlanRequest.schema_id
Schema to use, by id, or None to default to the sole in-window schema.
ReadPlanRequest.session_id
Limits the read to the topic’s data in this Session’s files, each over the part of its time span the Session holds; None adds no session restriction.
ReadPlanRequest.start_time
Inclusive window lower bound, absolute Unix-epoch nanoseconds.
May be None only when the request names a file_id, dataset_id or device_id; the bound then defaults to the earliest time of the topic’s data within the request’s restrictions, across every timeline source. An explicit bound narrows the read and never widens it.
ReadPlanRequest.timeline_source_id
Timeline source to resolve partition extents with, by id, or None.
ReadPlanRequest.timeline_source_name
Timeline source to resolve partition extents with, by name, or None.
ReadPlanScanTask
Bases: pydantic.BaseModel
One file to open, with the topic to read from it and the format and transformations needed to interpret it.
Which of the representations satisfying the governing selector backs a scan task is service policy and may change between releases; only the selector’s hard-filter matching rule is contract.
Parameters
data AnyAttributes
ReadPlanScanTask.format
The format the bytes are stored in; selects the decoder a consumer applies.
ReadPlanScanTask.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanScanTask.precedence
Where two scan tasks’ subtrees contain the same field, the read takes it from the higher-precedence one.
ReadPlanScanTask.subtree
The field subtree this scan task covers; None covers the whole schema.
ReadPlanScanTask.topic_name
Name of the topic whose records this scan task reads from its file.
A file can hold several topics (an MCAP file does when its channels carry different topic names), so a reader selects the topic’s records by this name. None when the plan does not name the topic; a reader then reads the file as holding a single topic.
ReadPlanScanTask.transformations
Transformations applied to produce this variant, in order; empty on the original.
ReadPlanSchemaRef
Bases: pydantic.BaseModel
Identifies the single topic schema the plan uses.
Parameters
data AnyAttributes
ReadPlanSchemaRef.checksum
Checksum of the schema’s content. A consumer can cache the schema by this value.
ReadPlanSchemaRef.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanTimestamp
Bases: pydantic.BaseModel
Where a partition’s row timestamps come from.
Timestamps are either read out of a schema field (kind is "schema_field", and field names which one) or taken from the storage envelope (message log or publish time), in which case no schema field is involved and field is None.
Parameters
data AnyAttributes
ReadPlanTimestamp.field
The schema field timestamps are read from; set exactly when kind is "schema_field".
ReadPlanTimestamp.kind
How timestamps are sourced: from a schema field, or from the storage envelope.
ReadPlanTimestamp.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
ReadPlanTimestamp.unit
Time unit of the designated field’s stored values (a TimeUnit value, e.g. "ms").
Only meaningful for a "schema_field" source, and only set when the schema declares the field’s unit. None when the schema does not record one; a consumer then treats non-self-describing values as nanoseconds, the unit the plan’s windows and offsets are recorded in. Envelope-derived timestamps (message log/publish time) are always assumed nanoseconds.
RepresentationOverride
Bases: pydantic.BaseModel
Applies a representation selector to one field subtree, overriding the request default.
Parameters
data AnyAttributes
RepresentationOverride.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
RepresentationOverride.selector
The selector to apply within that subtree.
RepresentationPreference
Bases: pydantic.BaseModel
Selects which stored variant of each field to read, per subtree.
A default selector applies to every field unless a more specific override covers it. Where several overrides cover a field, the one whose addressed subtree is the longest prefix of the field’s path wins; this rule is selector_for().
The governing selector and its matching rule are contract: a selector never substitutes a non-matching variant, and a read fails when a selector that sets any criterion is satisfied by no stored representation for a requested field — the plan never silently omits a field an explicit requirement covers. Which of the representations that satisfy the selector the service ultimately schedules is service policy and may change between releases.
Parameters
data AnyAttributes
RepresentationPreference.default
The selector applied to any field no override covers; matches anything when unset.
RepresentationPreference.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
RepresentationPreference.overrides
Per-subtree selector overrides, resolved longest-matching-prefix wins.
RepresentationPreference.selector_for()
Resolve the selector that governs the field at field_path, longest-matching-prefix wins.
An override applies when its addressed subtree path is a prefix of field_path; among applicable overrides the deepest subtree wins, and a field no override covers gets default.
Parameters
field_path tuple[str, .The path_in_schema components of the field whose selector is being resolved.
Returns
The governing selector.
Raises
ValueErrorAn override addresses its subtree by field_id. Resolving a field_id to a path takes the schema, which this value object does not hold; resolve every override address to its path form first.
RepresentationRecord
Bases: pydantic.BaseModel
One stored variant of a topic partition’s data, optionally narrowed to a subset of its fields.
A representation pairs a stored file with the data of a single topic partition. field_id narrows it to one field and the fields nested under it; None covers every field in the partition.
The same partition can have several representations that differ in storage_format, content_format, and transformations. A consumer picks the one whose attributes suit it: a viewer of image data, for example, may prefer a JPEG- or PNG-encoded variant over the untransformed original.
Parameters
data AnyAttributes
RepresentationRecord.content_format
The format of the data inside the stored file. For image data, this may be the image encoding (e.g. "jpeg", "png") on a transformed variant. None when unspecified.
RepresentationRecord.created
RepresentationRecord.created_by
Identity of the user who created this representation.
RepresentationRecord.field_id
The field this representation is narrowed to, covering that field and the fields nested under it. None when the representation covers every field in the partition.
RepresentationRecord.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
RepresentationRecord.modified
RepresentationRecord.modified_by
RepresentationRecord.representation_id
Unique identifier of this representation.
RepresentationRecord.size_bytes
Size in bytes of the file backing this representation, when known; None when the size is unavailable.
Populated on read-plan resolution so the plan can carry the backing file’s size onto its object refs. Write paths that upsert a representation leave it None.
RepresentationRecord.storage_format
Container the representation data is stored in (e.g. MCAP, Parquet).
RepresentationRecord.topic_part_id
Identifier of the topic partition this representation belongs to.
RepresentationRecord.transformations
The transformations applied to the source data to produce this variant, in the order applied. Empty on the untransformed original.
Each entry is a "<kind>:<param>" string whose <kind> is a TransformationKind member, e.g. ["downsample:0.5", "encode:jpeg"].
RepresentationSelector
Bases: pydantic.BaseModel
Selects which stored variant of a field to read when several are available.
A selector has three optional criteria — storage_format, content_format, and transformations — one for each attribute on which stored variants of the same field can differ (see RepresentationRecord). A criterion that is set is a requirement a variant must meet to be selected; a criterion left None places no requirement, and any value is acceptable.
A selector never falls back to a variant other than the one it describes. If any criterion is set and no stored variant of a requested field meets every requirement — whether the variants that exist all fall short, or the field has no stored variant at all — the read fails with an error rather than quietly leave out the field. Only under a selector with no criteria set is a field with no stored variant simply absent from the result; such a selector requires nothing, so nothing requested is missing.
Successor to roboto.domain.topics.RepresentationSelector, used by get_data().
Parameters
data AnyAttributes
RepresentationSelector.content_format
Required content encoding (e.g. "jpeg") by scalar equality; None does not constrain it.
There is no legacy carve-out: a representation whose content_format is None does not satisfy an explicit request.
RepresentationSelector.matches()
Return whether representation satisfies every set axis of this selector.
Parameters
representation RepresentationRecordReturn type
Attributes
RepresentationSelector.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
RepresentationSelector.raw()
Select the untransformed original (a representation with no transformations).
Return type
Attributes
RepresentationSelector.storage_format
Required container (e.g. MCAP, Parquet) by scalar equality; None does not constrain it.
RepresentationSelector.transformations
Required transformations; None does not constrain, () requires the untransformed original.
A non-empty tuple is all-of: every token must be satisfied by some descriptor on the representation, which may carry additional transformations. A token is either a bare kind (e.g. "downsample"), satisfied by any descriptor whose kind prefix equals it, or a full "<kind>:<param>" descriptor (e.g. "encode:jpeg"), satisfied only by an exact match. The grammar is open string matching; the recognized vocabulary (TransformationKind) is enforced by the service at the request boundary, which rejects a token naming an unrecognized kind.
SessionContext
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 AnyAttributes
SessionContext.end_time
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
SessionContext.start_time
Earliest time covered by the Session (Unix-epoch ns); the default start_time for get_data*. None when the Session includes no files.
SetTopicUnixOffsetRequest
Bases: pydantic.BaseModel
Request body for POST /v2/topics/id/<topic_id>/unix-offset.
Anchors one Session’s data on one topic to wall-clock time. See set_unix_offset() for the write’s reach.
Parameters
data AnyAttributes
SetTopicUnixOffsetRequest.model_config
model_config #Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].
SetTopicUnixOffsetRequest.session_id
Session whose data is anchored. Its files decide which of the topic’s stored data the write reaches; the topic’s data in files outside the Session keeps the anchor it already has.
SetTopicUnixOffsetRequest.unix_epoch_offset_ns
Wall-clock instant of stored time 0 for the Session’s data on this topic, in nanoseconds since the Unix epoch; each stored timestamp then reads as stored_time_ns + unix_epoch_offset_ns. Must fall after the Unix epoch, and must fit in the signed 64-bit integer the platform stores it in. Also accepts any roboto.time.Time at runtime, read as roboto.time.to_epoch_nanoseconds() reads it; convert with that function first to satisfy a type checker.
TIMESTAMP_FIELD_METADATA_KEY
Arrow field-metadata key marking the per-row timestamp column of a topic-data batch.
The timestamp column holds this key with the value b"true" (MARKER_FIELD_METADATA_VALUE); a column holding the key with any other value is a value column.
TimeWindow
Bases: pydantic.BaseModel
A closed time window in absolute nanoseconds since the Unix epoch; both bounds inclusive.
The same shape serves the window a plan resolves over and the window each partition’s rows are selected in.
Parameters
data AnyTopic
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
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.context Optional[TopicContext]Topic.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 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
ValueErrorThis 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()
Load an existing topic by its id.
Parameters
topic_id strIdentifier of the topic (ti_*).
roboto_client Optional[roboto.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()
Wrap an already-loaded topic identity record.
Parameters
The topic identity record to wrap.
roboto_client Optional[roboto.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()
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
start_time Optional[roboto.end_time Optional[roboto.fields_include Optional[collections.fields_exclude Optional[collections.prefer Optional[roboto.schema_id Optional[str]schema_checksum Optional[str]timeline_source_id Optional[str]timeline_source_name Optional[str]cache_policy roboto.cache_dir Union[str, pathlib.Yields
(timestamp, record) tuples for the in-window rows, filtered and projected per the arguments.
Raises
Return type
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()
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
start_time Optional[roboto.end_time Optional[roboto.fields_include Optional[collections.fields_exclude Optional[collections.prefer Optional[roboto.schema_id Optional[str]schema_checksum Optional[str]timeline_source_id Optional[str]timeline_source_name Optional[str]flatten boolExpand 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.cache_dir Union[str, pathlib.Returns
DataFrame of the in-window rows indexed by a timezone-aware DatetimeIndex.
Raises
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()
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.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.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.Field subtrees to project. None projects every field.
fields_exclude Optional[collections.Field subtrees to drop from the projection. None drops none.
prefer Optional[roboto.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.
cache_policy roboto.Whether fetched Parquet files are cached to local disk. MCAP data always streams.
cache_dir Union[str, pathlib.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_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.
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
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
Human-readable topic name (e.g. "/camera/image_raw"). Unique within an organization.
Topic.record
The underlying topic identity record.
Topic.set_unix_offset()
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:
- 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.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
ValueErrorThis 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().
TypeErroranchor is not one of the Time types.
ValueErroranchor 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
TopicContext
The scope a Topic reads within: one Session, file, dataset or device.
timestamp_column_index()
Locate the timestamp column: the field whose TIMESTAMP_FIELD_METADATA_KEY metadata is b"true".
The column is identified by metadata, never by name: a projected root field can legitimately carry any name, including the timestamp column’s conventional one.
Parameters
schema pyarrow.Raises
The schema does not contain exactly one marked column.
Return type