fujilib.sinks¶
sample_to_row() flattens a sample into one wide row whose columns
row_columns() fixes from the channels. A sink locks those columns before its
first row and writes samples; pipe() writes a recording to one (design §7.6).
fujilib.sinks.base ¶
Rows, schemas, the sink contract and :func:pipe (design §7.6).
Rows. :func:sample_to_row flattens a
:class:~fujilib.streaming.sample.Sample into one wide row, with columns
fixed by a set of channels:
- the header:
device,address,protocol,t_mono_ns,t_utc,t_midpoint_mono_ns,requested_at,received_at,latency_s; - per channel
N, the columns of :data:~fujilib.devices.models.READING_COLUMNSprefixedchN_(ch3_value,ch3_state, …); - the analyzer columns of :data:
~fujilib.devices.models.ANALYZER_COLUMNS; error_typeanderror_message,Noneon a successful poll.
Every value is float, int, str, bool or None; datetimes are
ISO 8601 strings. A successful row and an error row have exactly the same keys:
an error row carries None in every reading and analyzer column.
Schemas. A sink fixes its columns with a :class:SchemaLock before its
first row, from :func:row_columns for a set of channels, never from the
values, so a recording that starts with an error still types every reading
column. A sample with a channel the schema lacks is refused rather than
written without it.
Sinks follow :class:SampleSink: open, write_many, close, and
async with. :func:pipe writes a recording's batches to one.
BaseSink ¶
The lifecycle the bundled sinks share: open once, write, close once.
A subclass implements :meth:_open, :meth:_write (given rows under the
locked schema) and :meth:_close. Closing runs even when the caller is
cancelled.
A sink called name; its columns are locked now if channels are given.
Source code in src/fujilib/sinks/base.py
close
async
¶
Close the sink; again is a no-op. Completes even when the caller is cancelled.
Raises:
| Type | Description |
|---|---|
FujiSinkWriteError
|
the backing store failed to finish. |
Source code in src/fujilib/sinks/base.py
open
async
¶
Open the sink; again is a no-op.
Raises:
| Type | Description |
|---|---|
FujiSinkError
|
the sink was closed; a sink is used once. |
FujiSinkWriteError
|
the backing store cannot be opened. |
FujiSinkDependencyError
|
an optional library the sink needs is missing. |
Source code in src/fujilib/sinks/base.py
write_many
async
¶
Write samples; an empty sequence writes nothing.
Raises:
| Type | Description |
|---|---|
FujiSinkError
|
the sink is not open. |
FujiSinkSchemaError
|
a sample has a channel the locked columns lack. |
FujiSinkWriteError
|
the backing store refused the write. |
Source code in src/fujilib/sinks/base.py
SampleSink ¶
SchemaLock ¶
A sink's columns, fixed once from a set of channels.
Locked explicitly with :meth:lock, or by the first batch
:meth:rows sees, to the channels its samples carry (their union, in
channel order). Types always come from :func:row_columns.
A lock for the sink named sink, locked now if channels are given.
Source code in src/fujilib/sinks/base.py
lock ¶
Fix the columns to those of channels; locking again to the same is allowed.
Raises:
| Type | Description |
|---|---|
FujiSinkSchemaError
|
the lock holds other channels already. |
Source code in src/fujilib/sinks/base.py
rows ¶
Each sample's row; unlocked columns are locked to this batch's channels first.
Raises:
| Type | Description |
|---|---|
FujiSinkSchemaError
|
a sample has a channel the columns lack. |
Source code in src/fujilib/sinks/base.py
pipe
async
¶
Write every batch of a recording to the open sink, until the stream ends.
Samples are written in groups: once batch_size are waiting, and at the
latest flush_interval seconds after the first of them arrived, even if
no more come. When the pipe stops, by the stream ending, an error or
cancellation, it takes the batches already waiting in the stream and writes
everything it holds first (for up to 30 s, even when cancelled), unless the
sink itself failed. So a recording stopped with Ctrl-C has a row for every
poll its summary counts.
Returns:
| Type | Description |
|---|---|
AcquisitionSummary
|
What the pipe wrote: |
AcquisitionSummary
|
|
AcquisitionSummary
|
counters are in its |
Raises:
| Type | Description |
|---|---|
FujiValidationError
|
|
FujiSinkError
|
the sink failed. |
FujiSinkWriteError
|
the last samples could not be written in time. |
Source code in src/fujilib/sinks/base.py
row_columns ¶
The columns of every row for an analyzer with these established channels.
Reading and analyzer columns are nullable, because an error row carries
None in all of them.
Source code in src/fujilib/sinks/base.py
sample_channels ¶
The channels sample's row has columns for.
:attr:Sample.channels <fujilib.streaming.sample.Sample.channels>, or,
for a sample made without them, its frame's.
Source code in src/fujilib/sinks/base.py
sample_to_row ¶
Flatten sample into one wide row; see the module docstring for the layout.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
sample
|
Sample
|
The sample to flatten. |
required |
channels
|
Iterable[ChannelId] | None
|
The channels that fix the row's columns. When omitted, the
sample's own (:func: |
None
|
Source code in src/fujilib/sinks/base.py
fujilib.sinks._schema ¶
Column specifications for tabular sinks.
A sink fixes its schema before the first row from :func:fujilib.sinks.base.row_columns
rather than inferring types from the first batch, so a recording that starts
with an error row cannot lock every reading column as a string.
ColumnSpec
dataclass
¶
One column of a tabular schema.
Attributes:
| Name | Type | Description |
|---|---|---|
name |
str
|
Column name, verbatim from the row dict. |
python_type |
type[float] | type[int] | type[str] | type[bool]
|
The scalar type backing the column: :class: |
nullable |
bool
|
Whether the column may hold |
fujilib.sinks.memory ¶
A sink that keeps samples in memory: for tests, notebooks and short recordings.
InMemorySink ¶
Bases: BaseSink
Keeps every sample written, in order. It grows without bound.
Example::
async with InMemorySink() as sink, record(source, rate_hz=1, duration=10) as rec:
await pipe(rec, sink)
rows = sink.rows()
An empty sink; its columns are locked now if channels are given.
Source code in src/fujilib/sinks/memory.py
fujilib.sinks.csv ¶
A CSV file sink: one header line, then one line per sample.
- The header is the locked columns (:func:
~fujilib.sinks.base.row_columns), written when the columns are known: at :meth:CsvSink.openif the sink was given its channels, otherwise with the first batch. - Text is quoted and numbers are not, so
None(an empty field) and empty text ("") stay apart: Python's :mod:csvreader withquoting=csv.QUOTE_NOTNULLreads them back asNoneand"". That matters forchN_errors, which is""for "no errors" andNonewhen unknown. Booleans are the texttrue/false; floats are written so they read back exactly; the file is UTF-8. - The file is flushed after every write, so a process that dies loses at most the batch it was writing. An existing file is replaced.
- File I/O runs in a worker thread, never on the event loop.
CsvSink ¶
Bases: BaseSink
Write rows to a CSV file (see the module docstring).
A sink for path; its columns are locked now if channels are given.
Source code in src/fujilib/sinks/csv.py
csv_cell ¶
A row value as the CSV writer takes it: a boolean becomes true / false.
fujilib.sinks.parquet ¶
A Parquet file sink; needs the parquet extra (pip install 'fujilib[parquet]').
- The Arrow schema is the locked columns
(:func:
~fujilib.sinks.base.row_columns):float64,int64,boolandstring, nullable exactly where a row can holdNone. It is fixed when the columns are: at :meth:ParquetSink.openif the sink was given its channels, otherwise with the first batch. - Rows are gathered into row groups of
row_group_sizerows (1,000 by default), whatever the size of eachwrite_many, and the last, shorter group is written on close.pipe()writes about once a second, so one row group per write would make a day's recording at 1 Hz 86,400 row groups, whose metadatapyarrowkeeps in memory until the file is closed. Rows waiting for their group are held as rows and made into one Arrow table per row group: an Arrow table per write, of a row or two, would leave the process memory it never gets back, over a megabyte an hour at 1 Hz. - The file's key-value metadata carries
fujilib.versionand whatever the caller passes asmetadata. - A Parquet file is only readable once its footer is written, by
:meth:
~fujilib.sinks.base.BaseSink.close. Closing runs on cancellation and on Ctrl-C, even when the last rows cannot be written, but a process that is killed leaves an unreadable file; use the CSV sink where that matters. pyarrowis imported when the sink opens, so fujilib imports without it. Allpyarrowwork runs in a worker thread, never on the event loop.
Compression ¶
A Parquet compression codec.
ParquetSink ¶
ParquetSink(
path,
*,
channels=None,
compression="zstd",
row_group_size=DEFAULT_ROW_GROUP_SIZE,
metadata=None,
)
Bases: BaseSink
Write rows to a Parquet file (see the module docstring).
A sink for path; its columns are locked now if channels are given.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str | PathLike[str]
|
The file to write. |
required |
channels
|
Iterable[ChannelId] | None
|
The channels whose columns to write. |
None
|
compression
|
Compression
|
The codec. |
'zstd'
|
row_group_size
|
int
|
Rows per row group. |
DEFAULT_ROW_GROUP_SIZE
|
metadata
|
Mapping[str, str] | None
|
Key-value metadata for the file, text to text. |
None
|
Raises:
| Type | Description |
|---|---|
FujiValidationError
|
an unknown |
Source code in src/fujilib/sinks/parquet.py
require_pyarrow ¶
Check that pyarrow can be imported.
Raises:
| Type | Description |
|---|---|
FujiSinkDependencyError
|
it cannot; the |