Recording¶
record() polls one or more analyzers at a fixed rate and hands you one
batch per tick: a mapping from each analyzer's name to its Sample. A
sink writes the batches to memory, CSV or Parquet. Everything here is
read-only. The design is in design §7.6.
import anyio
from fujilib import ParquetSink, PollSourceAdapter, open_device, pipe, record
async def main() -> None:
async with await open_device(
"COM8", channel_map={"CH1": "co2", "CH2": "co", "CH3": "o2"}
) as anz:
source = PollSourceAdapter("zpa", anz)
channels = [c.channel for c in anz.channels]
async with (
ParquetSink("run.parquet", channels=channels) as sink,
record(source, rate_hz=1.0, duration=3600) as rec,
):
await pipe(rec, sink)
print(rec.summary)
anyio.run(main)
Or iterate the batches yourself:
async with record(source, rate_hz=1.0) as rec:
async for batch in rec:
sample = batch["zpa"]
if sample.error is None:
print(sample.frame.channel("CH3").value)
Samples and rows¶
One Sample per analyzer per tick carries the whole Frame: every
established channel and the analyzer's status, read in two Modbus
transactions. sample_to_row(sample) flattens it into one wide row of
scalars: a header (device, address, t_mono_ns, t_utc, latency_s,
…), then per channel chN_value, chN_raw, chN_decimals, chN_unit,
chN_gas, chN_label_source, chN_state, chN_valid, chN_hold,
chN_calibrating and chN_errors, then the analyzer's errors and alarms, then
error_type and error_message.
- Timestamps.
t_mono_nsandt_utcare the midpoint of the request and the reply of the block that holds every concentration. They are host times; the analyzer's own clock is never used for them. The analyzer's response time (15 s on the development unit) is applied inside the analyzer and is recorded in its metadata, not subtracted from timestamps. - The columns are fixed when the recording starts. Every sample carries
the channels established then (
Sample.channels), so every row of a recording, a failed poll included, has the same keys and the same types. A channel that is first seen alive later is left out of this recording's rows; assert every channel you need withchannel_map. - A failed poll is a row too. Its
frameisNone, itserroris set, and its row carriesNonein every reading and analyzer column, with the error's type and message. Gaps are recorded, never dropped. - Manual calibrations are not in the rows. A zero or span made at the
front panel marks its channels
calibrating. Feed the frames to aManualCalibrationTracker(fujilib.devices.panel) to get one event per calibration: its channels, how it ended, and the readings before and after. Frames read at 1 Hz can miss the second or two a calibration runs; the event then rests on an undocumented register, or says it is ambiguous.
from fujilib.devices.panel import ManualCalibrationTracker, PanelObservation
tracker = ManualCalibrationTracker()
for sample in samples: # of one analyzer, in order
if sample.frame is not None and (seen := PanelObservation.from_frame(sample.frame)):
if event := tracker.feed(seen):
print(event.kind, event.outcome, event.channels)
The schedule¶
Tick k is due k / rate_hz seconds after the first, which runs at once.
When a poll overruns so far that whole slots pass, the recorder skips them
and counts them in summary.samples_late; it never polls in a burst to catch
up. One analyzer's poll takes about 0.12 s through an FTDI adapter, so about
7-8 Hz is the practical ceiling for one station, and 1 Hz is the usual rate.
duration makes every tick due before it has passed, ceil(duration * rate_hz)
of them and at least one; without it, the recording runs until the async with
block exits.
When the consumer is slow¶
The stream holds buffer_size batches (64 by default). When it is full,
overflow decides:
OverflowPolicy |
What happens |
|---|---|
BLOCK (default) |
the recorder waits; the consumer sets the pace, and ticks missed meanwhile are counted late |
DROP_NEWEST |
the new batch is discarded |
DROP_OLDEST |
the oldest waiting batch is discarded to make room |
A batch is dropped whole, never split, and counted in
summary.samples_dropped.
Connection failures¶
A timeout, a damaged reply or any other failed poll is an error row, and the
recording goes on. A connection failure (the adapter unplugged, the port
gone) ends it: the tick's error row is still delivered, the stream ends, and
leaving the async with block raises FujiConnectionError.
To ride out a connection failure instead, pass a ReconnectPolicy:
Every tick of the outage is then an error row, and the analyzer is reopened on
the policy's schedule (0.5 s, 1, 2, 5, 10, then every 30 s by default).
Reopening opens the port again under the same settings and identifies the
analyzer, which must be the one that was open: the same serial number and type
code. Only an analyzer opened by port name can be reopened, not one opened on
a transport you passed in. You can also reopen one yourself with
await anz.reopen().
The rows after an outage are marked: for 90 s after the reopen, a reading that
would be ok has the state settling and is not valid, because the analyzer
may have been switched off and reports nothing of its own while it warms up
(see Measurement quality). The values are recorded
as read. open_device(..., settle_after_reopen_s=...) sets the period, and 0
turns it off.
The summary¶
rec.summary is updated live and finished when the recording stops, however
it stops. It counts polls, not rows.
| Field | Meaning |
|---|---|
samples_emitted |
batches put on the stream (and not dropped since); once stopped, exactly the batches the consumer received |
samples_late |
ticks skipped because a poll overran |
samples_dropped |
batches the overflow policy discarded, and batches still unread when the recording stopped |
error_samples |
samples of failed polls |
disconnects, reconnects |
connection failures (each outage ridden out, and the one that ends a recording), and outages ended by reopening |
max_drift_ms |
the latest a poll started after its slot |
target_total_samples |
the ticks a duration asked for |
For a recording that runs to its end,
samples_emitted + samples_dropped + samples_late == target_total_samples.
The traffic counters (requests, retries, failures by kind) are on
anz.session.counters.
Sinks¶
| Sink | Notes |
|---|---|
InMemorySink |
keeps every sample; for tests and short recordings |
CsvSink |
flushed after every write, so a process that dies loses only what was not written yet: what pipe() held (at most flush_interval seconds or batch_size samples) and what waited in the recording's buffer. Text is quoted and numbers are not, so an empty field is None and "" is empty text; read it back with csv.QUOTE_NOTNULL |
ParquetSink |
the parquet extra (pip install 'fujilib[parquet]'); zstd; rows gathered into row groups of 1,000 (row_group_size); the file's metadata carries fujilib.version and whatever you pass. Readable only once closed, which cancellation and Ctrl-C still do; a killed process leaves a file without its footer, whose complete row groups the repository's scripts/recover_parquet.py recovers |
Each fixes its columns from row_columns(channels) before the first row:
pass channels= to fix them at open(), or they are taken from the first
batch. Types never come from the values, so a recording that starts with an
error row still types ch3_value as a float. A sample with a channel the
columns lack is refused (FujiSinkSchemaError) rather than written without it.
File I/O runs in a worker thread, never on the event loop.
pipe(rec, sink, batch_size=64, flush_interval=1.0) writes the batches in
groups, at the latest one flush_interval after the first of a group arrived.
When it stops, by the recording ending, an error or cancellation, it takes the
batches already waiting and writes what it holds first. A file written by
pipe() has one row per analyzer for every batch the summary counts as
emitted, a recording stopped with Ctrl-C included.
Blocking code¶
fujilib.sync has the same recorder for code without an event loop:
from fujilib.sync import Fuji, PollSourceAdapter, SyncCsvSink, pipe, record
with (
Fuji.open("COM8", channel_map={"CH3": "o2"}) as anz,
record(PollSourceAdapter("zpa", anz), rate_hz=1.0, duration=60) as rec,
SyncCsvSink("run.csv", portal=anz.portal) as sink,
):
pipe(rec, sink)
for batch in rec: iterates the batches instead.
From the command line¶
fuji-stream prints each poll; fuji-capture records to a file with its
metadata beside it. See Commands.