Skip to content

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_COLUMNS prefixed chN_ (ch3_value, ch3_state, …);
  • the analyzer columns of :data:~fujilib.devices.models.ANALYZER_COLUMNS;
  • error_type and error_message, None on 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

BaseSink(name, channels=None)

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
def __init__(self, name: str, channels: Iterable[ChannelId] | None = None) -> None:
    """A sink called ``name``; its columns are locked now if ``channels`` are given."""
    self._name = name
    self._schema = SchemaLock(name, channels)
    self._state = "new"

is_open property

is_open

Whether the sink is open for writing.

schema property

schema

The sink's columns.

close async

close()

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
async def close(self) -> None:
    """Close the sink; again is a no-op. Completes even when the caller is cancelled.

    Raises:
        FujiSinkWriteError: the backing store failed to finish.
    """
    if self._state == "closed":
        return
    opened, self._state = self._state == "open", "closed"
    if opened:
        with anyio.CancelScope(shield=True):
            await self._close()

open async

open()

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
async def open(self) -> None:
    """Open the sink; again is a no-op.

    Raises:
        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.
    """
    if self._state == "open":
        return
    if self._state == "closed":
        msg = f"{self._name}: the sink was closed; open a new one"
        raise FujiSinkError(msg)
    await self._open()
    self._state = "open"

write_many async

write_many(samples)

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
async def write_many(self, samples: Sequence[Sample]) -> None:
    """Write ``samples``; an empty sequence writes nothing.

    Raises:
        FujiSinkError: the sink is not open.
        FujiSinkSchemaError: a sample has a channel the locked columns lack.
        FujiSinkWriteError: the backing store refused the write.
    """
    if self._state != "open":
        msg = f"{self._name}: write_many() needs an open sink"
        raise FujiSinkError(msg)
    if not samples:
        return
    await self._write(samples, self._schema.rows(samples))

SampleSink

Bases: Protocol

Where :func:pipe writes samples.

open and close are idempotent; async with opens and closes.

close async

close()

Finish writing and release what the sink holds.

Source code in src/fujilib/sinks/base.py
async def close(self) -> None:
    """Finish writing and release what the sink holds."""
    ...

open async

open()

Get ready to write.

Source code in src/fujilib/sinks/base.py
async def open(self) -> None:
    """Get ready to write."""
    ...

write_many async

write_many(samples)

Write samples, in order.

Source code in src/fujilib/sinks/base.py
async def write_many(self, samples: Sequence[Sample]) -> None:
    """Write ``samples``, in order."""
    ...

SchemaLock

SchemaLock(sink, channels=None)

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
def __init__(self, sink: str, channels: Iterable[ChannelId] | None = None) -> None:
    """A lock for the sink named ``sink``, locked now if ``channels`` are given."""
    self._sink = sink
    self._channels: tuple[ChannelId, ...] | None = None
    self._columns: tuple[ColumnSpec, ...] = ()
    if channels is not None:
        _ = self.lock(channels)

channels property

channels

The locked channels; empty before locking.

columns property

columns

The locked columns, in row order; empty before locking.

is_locked property

is_locked

Whether the columns are fixed.

lock

lock(channels)

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
def lock(self, channels: Iterable[ChannelId]) -> tuple[ColumnSpec, ...]:
    """Fix the columns to those of ``channels``; locking again to the same is allowed.

    Raises:
        FujiSinkSchemaError: the lock holds other channels already.
    """
    ordered = tuple(sorted(set(channels), key=lambda c: c.number))
    if self._channels is not None:
        if ordered != self._channels:
            msg = (
                f"{self._sink}: the columns are locked to {_names(self._channels)}, "
                f"not {_names(ordered)}"
            )
            raise FujiSinkSchemaError(msg)
        return self._columns
    self._channels = ordered
    self._columns = row_columns(ordered)
    return self._columns

rows

rows(samples)

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
def rows(self, samples: Sequence[Sample]) -> list[dict[str, Scalar]]:
    """Each sample's row; unlocked columns are locked to this batch's channels first.

    Raises:
        FujiSinkSchemaError: a sample has a channel the columns lack.
    """
    if self._channels is None:
        _ = self.lock(c for s in samples for c in sample_channels(s))
    locked = self.channels
    for sample in samples:
        extra = set(sample_channels(sample)) - set(locked)
        if extra:
            msg = (
                f"{self._sink}: the sample of {sample.device!r} has "
                f"{_names(sorted(extra, key=lambda c: c.number))}, which the columns, "
                f"locked to {_names(locked)}, do not"
            )
            raise FujiSinkSchemaError(msg)
    return [sample_to_row(s, locked) for s in samples]

pipe async

pipe(source, sink, *, batch_size=64, flush_interval=1.0)

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: samples_emitted counts batches (polls), and

AcquisitionSummary

error_samples the samples of failed polls. The recording's own

AcquisitionSummary

counters are in its summary.

Raises:

Type Description
FujiValidationError

source is not a recording or a memory object stream, or batch_size or flush_interval is invalid.

FujiSinkError

the sink failed.

FujiSinkWriteError

the last samples could not be written in time.

Source code in src/fujilib/sinks/base.py
async def pipe(
    source: Recording[Batch] | MemoryObjectReceiveStream[Batch],
    sink: SampleSink,
    *,
    batch_size: int = 64,
    flush_interval: float = 1.0,
) -> AcquisitionSummary:
    """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:
        What the pipe wrote: ``samples_emitted`` counts batches (polls), and
        ``error_samples`` the samples of failed polls. The recording's own
        counters are in its ``summary``.

    Raises:
        FujiValidationError: ``source`` is not a recording or a memory object
            stream, or ``batch_size`` or ``flush_interval`` is invalid.
        FujiSinkError: the sink failed.
        FujiSinkWriteError: the last samples could not be written in time.
    """
    stream = source.stream if isinstance(source, Recording) else source
    if not _is_memory_stream(stream):
        msg = f"source must be a Recording or its stream, got {type(source).__name__}"
        raise FujiValidationError(msg)
    if not _is_count(batch_size) or batch_size < 1:
        msg = f"batch_size must be an integer of at least 1, got {batch_size!r}"
        raise FujiValidationError(msg)
    if not _is_seconds(flush_interval) or flush_interval <= 0:
        msg = f"flush_interval must be a finite number of seconds above 0, got {flush_interval!r}"
        raise FujiValidationError(msg)
    writer = _PipeWriter(sink)
    try:
        await writer.run(stream, batch_size, flush_interval)
    finally:
        writer.drain(stream)
        await writer.finish()
    if writer.unwritten:
        msg = f"could not write the last {writer.unwritten} samples within {_FINAL_WRITE_S:g} s"
        raise FujiSinkWriteError(msg)
    return writer.summary

row_columns

row_columns(channels)

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
def row_columns(channels: Iterable[ChannelId]) -> tuple[ColumnSpec, ...]:
    """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.
    """
    specs = list(HEADER_COLUMNS)
    for channel in channels:
        prefix = _prefix(channel)
        specs.extend(
            ColumnSpec(prefix + c.name, c.python_type, nullable=True) for c in READING_COLUMNS
        )
    specs.extend(ColumnSpec(c.name, c.python_type, nullable=True) for c in ANALYZER_COLUMNS)
    specs.extend(_ERROR_COLUMNS)
    return tuple(specs)

sample_channels

sample_channels(sample)

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
def sample_channels(sample: Sample) -> tuple[ChannelId, ...]:
    """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.
    """
    if sample.channels:
        return sample.channels
    return sample.frame.channels if sample.frame is not None else ()

sample_to_row

sample_to_row(sample, channels=None)

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:sample_channels): a recorder sets them on every sample, failed polls included, so every row of a recording has the same keys.

None
Source code in src/fujilib/sinks/base.py
def sample_to_row(
    sample: Sample,
    channels: Iterable[ChannelId] | None = None,
) -> dict[str, Scalar]:
    """Flatten ``sample`` into one wide row; see the module docstring for the layout.

    Args:
        sample: The sample to flatten.
        channels: The channels that fix the row's columns. When omitted, the
            sample's own (:func:`sample_channels`): a recorder sets them on
            every sample, failed polls included, so every row of a recording
            has the same keys.
    """
    frame = sample.frame
    established = tuple(channels) if channels is not None else sample_channels(sample)
    row: dict[str, Scalar] = {
        "device": sample.device,
        "address": sample.address,
        "protocol": sample.protocol.value,
        "t_mono_ns": sample.t_mono_ns,
        "t_utc": sample.t_utc.isoformat(),
        "t_midpoint_mono_ns": sample.t_midpoint_mono_ns,
        "requested_at": sample.requested_at.isoformat(),
        "received_at": sample.received_at.isoformat(),
        "latency_s": sample.latency_s,
    }
    readings = {r.channel: r for r in frame.readings} if frame is not None else {}
    for channel in established:
        prefix = _prefix(channel)
        reading = readings.get(channel)
        values = reading.as_dict() if reading is not None else {}
        row.update({prefix + c.name: values.get(c.name) for c in READING_COLUMNS})
    analyzer = frame.analyzer.as_dict() if frame is not None and frame.analyzer is not None else {}
    row.update({c.name: analyzer.get(c.name) for c in ANALYZER_COLUMNS})
    error = sample.error
    if error is not None:
        cls = type(error)
        row["error_type"] = f"{cls.__module__}.{cls.__qualname__}"
        row["error_message"] = str(error)
    else:
        row["error_type"] = None
        row["error_message"] = None
    return row

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

ColumnSpec(name, python_type, nullable)

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:float, :class:int, :class:str or :class:bool.

nullable bool

Whether the column may hold None.

fujilib.sinks.memory

A sink that keeps samples in memory: for tests, notebooks and short recordings.

InMemorySink

InMemorySink(*, channels=None)

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
def __init__(self, *, channels: Iterable[ChannelId] | None = None) -> None:
    """An empty sink; its columns are locked now if ``channels`` are given."""
    super().__init__("memory", channels)
    self._samples: list[Sample] = []
    self._rows: list[dict[str, Scalar]] = []

samples property

samples

The samples written, in order. Kept after :meth:close.

rows

rows()

The rows written, under the sink's locked columns.

Source code in src/fujilib/sinks/memory.py
def rows(self) -> list[dict[str, Scalar]]:
    """The rows written, under the sink's locked columns."""
    return list(self._rows)

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.open if 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:csv reader with quoting=csv.QUOTE_NOTNULL reads them back as None and "". That matters for chN_errors, which is "" for "no errors" and None when unknown. Booleans are the text true / 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

CsvSink(path, *, channels=None)

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
def __init__(
    self, path: str | PathLike[str], *, channels: Iterable[ChannelId] | None = None
) -> None:
    """A sink for ``path``; its columns are locked now if ``channels`` are given."""
    super().__init__(f"csv:{path}", channels)
    self._path = Path(path)
    self._file: IO[str] | None = None
    self._header_written = False

path property

path

The file written.

csv_cell

csv_cell(value)

A row value as the CSV writer takes it: a boolean becomes true / false.

Source code in src/fujilib/sinks/csv.py
def csv_cell(value: Scalar) -> str | int | float | None:
    """A row value as the CSV writer takes it: a boolean becomes ``true`` / ``false``."""
    if isinstance(value, bool):
        return "true" if value else "false"
    return value

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, bool and string, nullable exactly where a row can hold None. It is fixed when the columns are: at :meth:ParquetSink.open if the sink was given its channels, otherwise with the first batch.
  • Rows are gathered into row groups of row_group_size rows (1,000 by default), whatever the size of each write_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 metadata pyarrow keeps 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.version and whatever the caller passes as metadata.
  • 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.
  • pyarrow is imported when the sink opens, so fujilib imports without it. All pyarrow work runs in a worker thread, never on the event loop.

Compression

Compression = Literal[
    "zstd", "snappy", "gzip", "brotli", "lz4", "none"
]

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 compression, a row_group_size below 1, or metadata that is not text.

Source code in src/fujilib/sinks/parquet.py
def __init__(
    self,
    path: str | PathLike[str],
    *,
    channels: Iterable[ChannelId] | None = None,
    compression: Compression = "zstd",
    row_group_size: int = DEFAULT_ROW_GROUP_SIZE,
    metadata: Mapping[str, str] | None = None,
) -> None:
    """A sink for ``path``; its columns are locked now if ``channels`` are given.

    Args:
        path: The file to write.
        channels: The channels whose columns to write.
        compression: The codec.
        row_group_size: Rows per row group.
        metadata: Key-value metadata for the file, text to text.

    Raises:
        FujiValidationError: an unknown ``compression``, a ``row_group_size``
            below 1, or ``metadata`` that is not text.
    """
    if compression not in COMPRESSIONS:
        known = ", ".join(sorted(COMPRESSIONS))
        msg = f"compression must be one of {known}, got {compression!r}"
        raise FujiValidationError(msg)
    if not _is_count(row_group_size) or row_group_size < 1:
        msg = f"row_group_size must be an integer of at least 1, got {row_group_size!r}"
        raise FujiValidationError(msg)
    extra = dict(metadata or {})
    if not all(_is_text(k) and _is_text(v) for k, v in extra.items()):
        msg = "metadata must map text to text"
        raise FujiValidationError(msg)
    super().__init__(f"parquet:{path}", channels)
    self._path = Path(path)
    self._compression: Compression = compression
    self._row_group_size = row_group_size
    self._metadata: Mapping[str, str] = MappingProxyType(
        {"fujilib.version": __version__, **extra}
    )
    self._writer: pq.ParquetWriter | None = None
    self._arrow_schema: pa.Schema | None = None
    self._waiting: list[dict[str, Scalar]] = []

metadata property

metadata

The file's key-value metadata.

path property

path

The file written.

require_pyarrow

require_pyarrow()

Check that pyarrow can be imported.

Raises:

Type Description
FujiSinkDependencyError

it cannot; the parquet extra is missing.

Source code in src/fujilib/sinks/parquet.py
def require_pyarrow() -> None:
    """Check that ``pyarrow`` can be imported.

    Raises:
        FujiSinkDependencyError: it cannot; the ``parquet`` extra is missing.
    """
    try:
        _ = importlib.import_module("pyarrow")
        _ = importlib.import_module("pyarrow.parquet")
    except ImportError as exc:
        msg = f"the Parquet sink needs pyarrow, which is not installed; {_INSTALL_HINT}"
        raise FujiSinkDependencyError(msg) from exc