How to Stream Data¶
Stream data to Zelos at any rate: steady, in bursts, from many sources, or with your own timestamps.
The Python samples are complete scripts. The Rust and Go samples are the body of main. They assume the router, the publish client, and a source set up as in the Quickstart.
Basic Streaming¶
The simplest way to stream data is to log events in a loop:
import time
import random # For demo purposes
import zelos_sdk
zelos_sdk.init()
source = zelos_sdk.TraceSource("sensors")
# Define schema (optional but recommended)
temperature_event = source.add_event("temperature", [
zelos_sdk.TraceEventFieldMetadata("value", zelos_sdk.DataType.Float64, "°C")
])
while True:
# Read your sensor (simulated here)
temperature = 20.0 + random.uniform(-2, 2)
# Log to Zelos
temperature_event.log(value=temperature)
# Control rate
time.sleep(1.0) # 1 Hz
use std::time::{Duration, Instant};
// Define schema (required in Rust)
let temp_event = source
.build_event("temperature")
.add_f64_field("value", Some("°C".to_string()))
.build()?;
let mut interval = tokio::time::interval(Duration::from_secs(1)); // 1 Hz
let start = Instant::now();
loop {
interval.tick().await;
// Read your sensor (simulated here)
let temperature = 20.0 + 2.0 * (start.elapsed().as_secs_f64() * 0.1).sin();
temp_event
.build()
.try_insert_f64("value", temperature)?
.emit()?;
}
// Define schema (required in Go)
unitC := "°C"
tempEvent, err := source.BuildEvent("temperature").
AddFloat64Field("value", &unitC).
Build()
if err != nil {
log.Fatal(err)
}
ticker := time.NewTicker(time.Second) // 1 Hz
defer ticker.Stop()
for range ticker.C {
// Read your sensor (simulated here)
temperature := 20.0 + rand.Float64()*4 - 2
builder, _ := tempEvent.Build().TryInsertFloat64("value", temperature)
if err := builder.Emit(); err != nil {
log.Printf("Emit failed: %v", err)
}
}
High-Frequency Streaming¶
Hold a Steady Rate¶
Compensate for loop overhead so the rate does not drift:
import time
import math
import zelos_sdk
zelos_sdk.init()
source = zelos_sdk.TraceSource("high_freq")
data_event = source.add_event("data", [
zelos_sdk.TraceEventFieldMetadata("value", zelos_sdk.DataType.Float64, "V")
])
# 1 kHz = 1 ms period
period = 0.001
start_time = time.time()
next_time = start_time
while True:
current_time = time.time() - start_time
# Generate signal (example: 100 Hz sine wave)
value = math.sin(2 * math.pi * 100 * current_time)
data_event.log(value=value)
# Sleep until the next period, not for a fixed time
next_time += period
sleep_duration = next_time - time.time()
if sleep_duration > 0:
time.sleep(sleep_duration)
use std::time::{Duration, Instant};
use tokio::time::MissedTickBehavior;
let data_event = source
.build_event("data")
.add_f64_field("value", Some("V".to_string()))
.build()?;
// 1 kHz = 1 ms period
let mut interval = tokio::time::interval(Duration::from_millis(1));
// Skip missed ticks instead of bursting
interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
let start = Instant::now();
loop {
interval.tick().await;
let elapsed = start.elapsed().as_secs_f64();
// Generate signal (example: 100 Hz sine wave)
let value = (2.0 * std::f64::consts::PI * 100.0 * elapsed).sin();
// Waits for channel space without blocking the runtime thread
data_event
.build()
.try_insert_f64("value", value)?
.emit_async()
.await?;
}
emit() blocks the calling thread while the router channel is full. Inside async code, use emit_async() (or emit_at_async(time_ns)) so the wait does not hold a runtime thread.
unitVolt := "V"
dataEvent, err := source.BuildEvent("data").
AddFloat64Field("value", &unitVolt).
Build()
if err != nil {
log.Fatal(err)
}
// 1 kHz = 1 ms period
ticker := time.NewTicker(time.Millisecond)
defer ticker.Stop()
start := time.Now()
for range ticker.C {
elapsed := time.Since(start).Seconds()
// Generate signal (example: 100 Hz sine wave)
value := math.Sin(2 * math.Pi * 100 * elapsed)
builder, _ := dataEvent.Build().TryInsertFloat64("value", value)
if err := builder.Emit(); err != nil {
log.Printf("Dropped sample: %v", err)
}
}
Emit does not wait. When the router channel is full, it drops the event and returns an error. Check that error at high rates.
Log Many Events per Call¶
Python only
log_many and log_batch are part of the Python SDK. The Rust and Go SDKs log one event per call.
Each log call crosses from Python into the SDK's native core. When your hardware hands you a block of samples at once, such as a DAQ buffer or a set of decoded CAN frames, log the whole block in one call.
TraceSource.log_many takes a list of (time_ns, event_name, fields) tuples. Each entry carries its own timestamp, so samples keep their real spacing. One call can mix several event names.
import math
import time
import zelos_sdk
zelos_sdk.init()
source = zelos_sdk.TraceSource("daq")
source.add_event("voltage", [
zelos_sdk.TraceEventFieldMetadata("value", zelos_sdk.DataType.Float64, "V")
])
SAMPLE_PERIOD_NS = 100_000 # 10 kHz
BLOCK_SIZE = 100 # samples per read
def read_block():
"""Return the first sample's time and a block of samples (simulated here)."""
t0 = time.time_ns()
samples = [
math.sin(2 * math.pi * 100 * (t0 + i * SAMPLE_PERIOD_NS) / 1e9)
for i in range(BLOCK_SIZE)
]
return t0, samples
while True:
t0, samples = read_block()
source.log_many([
(t0 + i * SAMPLE_PERIOD_NS, "voltage", {"value": value})
for i, value in enumerate(samples)
])
time.sleep(BLOCK_SIZE * SAMPLE_PERIOD_NS / 1e9)
- Pass a list of tuples.
log_manyrejects a generator and rejects entries that are lists. - An event name with no schema yet gets one from its first dict, as with
source.log. - Use the schema's field names.
log_manyskips a key that is not in the schema instead of raising, so a typo loses that value without an error. - Use flat dicts.
log_manydoes not log nested dicts. Noneis not a valid field value here.log_manyraisesValueErrorfor it.
Log Columnar Data¶
When your data is already in columns, such as NumPy arrays or a pandas DataFrame, pass it to TraceSource.log_batch(event_name, data) as a pyarrow.RecordBatch or pyarrow.Table. log_batch sends the columns without converting each value in Python.
import math
import time
import pyarrow as pa
import zelos_sdk
zelos_sdk.init()
source = zelos_sdk.TraceSource("daq")
SAMPLE_PERIOD_NS = 100_000 # 10 kHz
BLOCK_SIZE = 1_000 # samples per read
while True:
# Read a block of samples (simulated here)
t0 = time.time_ns()
time_ns = [t0 + i * SAMPLE_PERIOD_NS for i in range(BLOCK_SIZE)]
ch0 = [math.sin(2 * math.pi * 100 * t / 1e9) for t in time_ns]
ch1 = [math.cos(2 * math.pi * 100 * t / 1e9) for t in time_ns]
batch = pa.record_batch({
"time_ns": pa.array(time_ns, type=pa.timestamp("ns")),
"ch0": pa.array(ch0, type=pa.float64()),
"ch1": pa.array(ch1, type=pa.float64()),
})
source.log_batch("channels", batch)
time.sleep(BLOCK_SIZE * SAMPLE_PERIOD_NS / 1e9)
pa.array also accepts NumPy arrays. For a pandas DataFrame, build the table with pa.Table.from_pandas(df, preserve_index=False).
- The first column is the timestamp. Its Arrow type must be
timestamp("ns"), with or without a time zone.log_batchreads the first column as time whatever its name, and does not check its type. A batch with another time type, such asint64ortimestamp("us"), produces a trace without usable times. - Every other column becomes a field, named after the column. Supported column types are signed and unsigned integers,
float32,float64,bool,string,binary, andtimestamp("ns"). log_batchdrops columns of any other Arrow type without an error. pandas text columns becomelarge_string, so cast them tostringfirst.- The first batch for an event name registers its schema, without units. To give fields units, register the event with
add_eventfirst.log_batchthen uses that schema, so the columns must match its field names and types. - Keep the same columns in every batch for one event name. Later batches are not checked against the schema.
log_batchsorts each batch by time.
Burst Streaming¶
Stream data in bursts when events occur:
import time
import zelos_sdk
zelos_sdk.init()
source = zelos_sdk.TraceSource("events")
# Define event schemas
source.add_event("event_start", [
zelos_sdk.TraceEventFieldMetadata("trigger", zelos_sdk.DataType.String),
zelos_sdk.TraceEventFieldMetadata("timestamp_ms", zelos_sdk.DataType.Int64)
])
source.add_event("event_sample", [
zelos_sdk.TraceEventFieldMetadata("value", zelos_sdk.DataType.Float64),
zelos_sdk.TraceEventFieldMetadata("index", zelos_sdk.DataType.UInt32)
])
source.add_event("event_end", [
zelos_sdk.TraceEventFieldMetadata("duration_ms", zelos_sdk.DataType.Float64),
zelos_sdk.TraceEventFieldMetadata("sample_count", zelos_sdk.DataType.UInt32)
])
def handle_burst_event(trigger_name: str, samples: list):
"""Log a burst of data when an event occurs"""
start_time = time.time()
# Log event start
source.log("event_start", {
"trigger": trigger_name,
"timestamp_ms": int(start_time * 1000)
})
# Log burst of samples
for i, value in enumerate(samples):
source.log("event_sample", {
"value": value,
"index": i
})
# Log event end
duration_ms = (time.time() - start_time) * 1000
source.log("event_end", {
"duration_ms": duration_ms,
"sample_count": len(samples)
})
while True:
# Wait for trigger condition
time.sleep(5)
# Generate burst data (simulated)
burst_data = [i * 0.1 for i in range(100)]
handle_burst_event("threshold_exceeded", burst_data)
To log a large burst in one call, use log_many.
use std::error::Error;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use zelos::trace::source::TraceSourceEvent;
// Define event schemas
let start_event = source
.build_event("event_start")
.add_string_field("trigger", None)
.add_i64_field("timestamp_ms", None)
.build()?;
let sample_event = source
.build_event("event_sample")
.add_f64_field("value", None)
.add_u32_field("index", None)
.build()?;
let end_event = source
.build_event("event_end")
.add_f64_field("duration_ms", None)
.add_u32_field("sample_count", None)
.build()?;
fn handle_burst_event(
trigger_name: &str,
samples: &[f64],
start_event: &TraceSourceEvent,
sample_event: &TraceSourceEvent,
end_event: &TraceSourceEvent,
) -> Result<(), Box<dyn Error>> {
let start_time = Instant::now();
// Log event start
start_event
.build()
.try_insert_string("trigger", trigger_name.to_string())?
.try_insert_i64(
"timestamp_ms",
SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as i64,
)?
.emit()?;
// Log burst of samples
for (i, value) in samples.iter().enumerate() {
sample_event
.build()
.try_insert_f64("value", *value)?
.try_insert_u32("index", i as u32)?
.emit()?;
}
// Log event end
let duration_ms = start_time.elapsed().as_secs_f64() * 1000.0;
end_event
.build()
.try_insert_f64("duration_ms", duration_ms)?
.try_insert_u32("sample_count", samples.len() as u32)?
.emit()?;
Ok(())
}
loop {
// Wait for trigger condition
tokio::time::sleep(Duration::from_secs(5)).await;
// Generate burst data (simulated)
let burst_data: Vec<f64> = (0..100).map(|i| i as f64 * 0.1).collect();
handle_burst_event(
"threshold_exceeded",
&burst_data,
&start_event,
&sample_event,
&end_event,
)?;
}
// Define event schemas
startEvent, err := source.BuildEvent("event_start").
AddStringField("trigger", nil).
AddInt64Field("timestamp_ms", nil).
Build()
if err != nil {
log.Fatal(err)
}
sampleEvent, err := source.BuildEvent("event_sample").
AddFloat64Field("value", nil).
AddUint32Field("index", nil).
Build()
if err != nil {
log.Fatal(err)
}
endEvent, err := source.BuildEvent("event_end").
AddFloat64Field("duration_ms", nil).
AddUint32Field("sample_count", nil).
Build()
if err != nil {
log.Fatal(err)
}
handleBurstEvent := func(triggerName string, samples []float64) error {
startTime := time.Now()
// Log event start
builder, _ := startEvent.Build().TryInsertString("trigger", triggerName)
builder, _ = builder.TryInsertInt64("timestamp_ms", startTime.UnixMilli())
if err := builder.Emit(); err != nil {
return err
}
// Log burst of samples
for i, value := range samples {
b, _ := sampleEvent.Build().TryInsertFloat64("value", value)
b, _ = b.TryInsertUint32("index", uint32(i))
if err := b.Emit(); err != nil {
return err
}
}
// Log event end
durationMs := time.Since(startTime).Seconds() * 1000
builder, _ = endEvent.Build().TryInsertFloat64("duration_ms", durationMs)
builder, _ = builder.TryInsertUint32("sample_count", uint32(len(samples)))
return builder.Emit()
}
for {
// Wait for trigger condition
time.Sleep(5 * time.Second)
// Generate burst data (simulated)
burstData := make([]float64, 100)
for i := range burstData {
burstData[i] = float64(i) * 0.1
}
if err := handleBurstEvent("threshold_exceeded", burstData); err != nil {
log.Printf("Error handling burst: %v", err)
}
}
Emit returns an error and drops the event when the router channel is full. Keep a burst within the channel size, or retry on error.
Multiple Sources in Parallel¶
Stream from different components at different rates:
import threading
import time
import random
import zelos_sdk
zelos_sdk.init()
def stream_motor():
"""100 Hz motor telemetry"""
source = zelos_sdk.TraceSource("motor")
telemetry = source.add_event("telemetry", [
zelos_sdk.TraceEventFieldMetadata("rpm", zelos_sdk.DataType.Float64, "rpm"),
zelos_sdk.TraceEventFieldMetadata("torque", zelos_sdk.DataType.Float64, "Nm")
])
while True:
telemetry.log(
rpm=2000 + random.uniform(-100, 100),
torque=50 + random.uniform(-5, 5)
)
time.sleep(0.01) # 100 Hz
def stream_battery():
"""1 Hz battery status"""
source = zelos_sdk.TraceSource("battery")
status = source.add_event("status", [
zelos_sdk.TraceEventFieldMetadata("voltage", zelos_sdk.DataType.Float64, "V"),
zelos_sdk.TraceEventFieldMetadata("current", zelos_sdk.DataType.Float64, "A"),
zelos_sdk.TraceEventFieldMetadata("soc", zelos_sdk.DataType.Float64, "%")
])
soc = 85.0 # State of charge
while True:
# Simulate battery discharge
soc = max(20.0, soc - 0.1)
status.log(
voltage=48.0 + random.uniform(-0.5, 0.5),
current=random.uniform(-10, 50),
soc=soc
)
time.sleep(1.0) # 1 Hz
def stream_gps():
"""10 Hz GPS updates"""
source = zelos_sdk.TraceSource("gps")
position = source.add_event("position", [
zelos_sdk.TraceEventFieldMetadata("lat", zelos_sdk.DataType.Float64, "deg"),
zelos_sdk.TraceEventFieldMetadata("lon", zelos_sdk.DataType.Float64, "deg"),
zelos_sdk.TraceEventFieldMetadata("alt", zelos_sdk.DataType.Float64, "m")
])
# Base coordinates (Palo Alto, CA)
base_lat, base_lon = 37.4419, -122.1430
while True:
position.log(
lat=base_lat + random.uniform(-0.001, 0.001),
lon=base_lon + random.uniform(-0.001, 0.001),
alt=30 + random.uniform(-1, 1)
)
time.sleep(0.1) # 10 Hz
# Run all streams in parallel
threads = [
threading.Thread(target=stream_motor, daemon=True),
threading.Thread(target=stream_battery, daemon=True),
threading.Thread(target=stream_gps, daemon=True)
]
for t in threads:
t.start()
# Keep main thread alive
try:
for t in threads:
t.join()
except KeyboardInterrupt:
print("Stopping streams...")
This sample is a whole file. It replaces main from the Quickstart.
use std::error::Error;
use std::time::{Duration, Instant};
use zelos::TraceSource;
async fn stream_motor(source: TraceSource) -> Result<(), Box<dyn Error>> {
let telemetry = source
.build_event("telemetry")
.add_f64_field("rpm", Some("rpm".to_string()))
.add_f64_field("torque", Some("Nm".to_string()))
.build()?;
let mut interval = tokio::time::interval(Duration::from_millis(10)); // 100 Hz
let start = Instant::now();
loop {
interval.tick().await;
let t = start.elapsed().as_secs_f64();
telemetry
.build()
.try_insert_f64("rpm", 2000.0 + 100.0 * t.sin())?
.try_insert_f64("torque", 50.0 + 5.0 * t.cos())?
.emit_async()
.await?;
}
}
async fn stream_battery(source: TraceSource) -> Result<(), Box<dyn Error>> {
let status = source
.build_event("status")
.add_f64_field("voltage", Some("V".to_string()))
.add_f64_field("current", Some("A".to_string()))
.add_f64_field("soc", Some("%".to_string()))
.build()?;
let mut interval = tokio::time::interval(Duration::from_secs(1)); // 1 Hz
let start = Instant::now();
let mut soc: f64 = 85.0; // State of charge
loop {
interval.tick().await;
let t = start.elapsed().as_secs_f64();
// Simulate battery discharge
soc = (soc - 0.1).max(20.0);
status
.build()
.try_insert_f64("voltage", 48.0 + 0.5 * t.sin())?
.try_insert_f64("current", 20.0 + 30.0 * t.cos())?
.try_insert_f64("soc", soc)?
.emit_async()
.await?;
}
}
async fn stream_gps(source: TraceSource) -> Result<(), Box<dyn Error>> {
let position = source
.build_event("position")
.add_f64_field("lat", Some("deg".to_string()))
.add_f64_field("lon", Some("deg".to_string()))
.add_f64_field("alt", Some("m".to_string()))
.build()?;
let mut interval = tokio::time::interval(Duration::from_millis(100)); // 10 Hz
let start = Instant::now();
// Base coordinates (Palo Alto, CA)
let base_lat = 37.4419;
let base_lon = -122.1430;
loop {
interval.tick().await;
let t = start.elapsed().as_secs_f64();
position
.build()
.try_insert_f64("lat", base_lat + 0.001 * t.sin())?
.try_insert_f64("lon", base_lon + 0.001 * t.cos())?
.try_insert_f64("alt", 30.0 + t.sin())?
.emit_async()
.await?;
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
// Set up the router and publish client as in the Quickstart
// ...
// Create sources
let motor = TraceSource::new("motor", router.sender());
let battery = TraceSource::new("battery", router.sender());
let gps = TraceSource::new("gps", router.sender());
// Run all streams concurrently; the first error stops them all
tokio::try_join!(
stream_motor(motor),
stream_battery(battery),
stream_gps(gps),
)?;
Ok(())
}
This sample is a whole file. It replaces main from the Quickstart.
package main
import (
"log"
"math"
"math/rand"
"sync"
"time"
zelos "github.com/zeloscloud/zelos/go"
)
func streamMotor(source *zelos.TraceSource, wg *sync.WaitGroup) {
defer wg.Done()
unitRpm := "rpm"
unitNm := "Nm"
telemetry, err := source.BuildEvent("telemetry").
AddFloat64Field("rpm", &unitRpm).
AddFloat64Field("torque", &unitNm).
Build()
if err != nil {
log.Fatal(err)
}
ticker := time.NewTicker(10 * time.Millisecond) // 100 Hz
defer ticker.Stop()
for range ticker.C {
builder, _ := telemetry.Build().
TryInsertFloat64("rpm", 2000+rand.Float64()*200-100)
builder, _ = builder.
TryInsertFloat64("torque", 50+rand.Float64()*10-5)
if err := builder.Emit(); err != nil {
log.Printf("motor: %v", err)
}
}
}
func streamBattery(source *zelos.TraceSource, wg *sync.WaitGroup) {
defer wg.Done()
unitV := "V"
unitA := "A"
unitPct := "%"
status, err := source.BuildEvent("status").
AddFloat64Field("voltage", &unitV).
AddFloat64Field("current", &unitA).
AddFloat64Field("soc", &unitPct).
Build()
if err != nil {
log.Fatal(err)
}
ticker := time.NewTicker(time.Second) // 1 Hz
defer ticker.Stop()
soc := 85.0 // State of charge
for range ticker.C {
// Simulate battery discharge
soc = math.Max(20.0, soc-0.1)
builder, _ := status.Build().
TryInsertFloat64("voltage", 48.0+rand.Float64()-0.5)
builder, _ = builder.
TryInsertFloat64("current", rand.Float64()*60-10)
builder, _ = builder.
TryInsertFloat64("soc", soc)
if err := builder.Emit(); err != nil {
log.Printf("battery: %v", err)
}
}
}
func streamGPS(source *zelos.TraceSource, wg *sync.WaitGroup) {
defer wg.Done()
unitDeg := "deg"
unitM := "m"
position, err := source.BuildEvent("position").
AddFloat64Field("lat", &unitDeg).
AddFloat64Field("lon", &unitDeg).
AddFloat64Field("alt", &unitM).
Build()
if err != nil {
log.Fatal(err)
}
ticker := time.NewTicker(100 * time.Millisecond) // 10 Hz
defer ticker.Stop()
// Base coordinates (Palo Alto, CA)
baseLat := 37.4419
baseLon := -122.1430
for range ticker.C {
builder, _ := position.Build().
TryInsertFloat64("lat", baseLat+rand.Float64()*0.002-0.001)
builder, _ = builder.
TryInsertFloat64("lon", baseLon+rand.Float64()*0.002-0.001)
builder, _ = builder.
TryInsertFloat64("alt", 30+rand.Float64()*2-1)
if err := builder.Emit(); err != nil {
log.Printf("gps: %v", err)
}
}
}
func main() {
// Set up the router and publish client as in the Quickstart
// ...
// Create sources
motor, err := zelos.NewTraceSource("motor", sender)
if err != nil {
log.Fatal(err)
}
battery, err := zelos.NewTraceSource("battery", sender)
if err != nil {
log.Fatal(err)
}
gps, err := zelos.NewTraceSource("gps", sender)
if err != nil {
log.Fatal(err)
}
defer motor.Close()
defer battery.Close()
defer gps.Close()
// Run all streams in parallel
var wg sync.WaitGroup
wg.Add(3)
go streamMotor(motor, &wg)
go streamBattery(battery, &wg)
go streamGPS(gps, &wg)
wg.Wait()
}
Stream with Custom Timestamps¶
Use specific timestamps for replay or synchronization. Timestamps are nanoseconds since the Unix epoch.
import time
import zelos_sdk
zelos_sdk.init()
source = zelos_sdk.TraceSource("replay")
measurement = source.add_event("measurement", [
zelos_sdk.TraceEventFieldMetadata("value", zelos_sdk.DataType.Float64)
])
# Method 1: Current timestamp (automatic)
measurement.log(value=1.0)
# Method 2: Specific timestamp using log_at
specific_time_ns = 1699564234567890123
measurement.log_at(specific_time_ns, value=2.0)
# Method 3: Calculated timestamp
past_time_ns = time.time_ns() - (60 * 1_000_000_000) # 1 minute ago
measurement.log_at(past_time_ns, value=3.0)
# Method 4: Synchronized timestamps
def stream_synchronized(sync_offset_ns: int):
"""Stream with timestamps on another clock"""
while True:
timestamp_ns = time.time_ns() + sync_offset_ns
value = read_sensor() # Your sensor reading function
measurement.log_at(timestamp_ns, value=value)
time.sleep(0.1)
# Method 5: Replay historical data with its original timestamps
def replay_historical_data(records: list):
source.log_many([
(record["timestamp_ns"], "measurement", {"value": record["value"]})
for record in records
])
Without an event handle, source.log_at(time_ns, "measurement", {"value": 2.0}) does the same as Method 2.
use std::error::Error;
use std::time::Duration;
use zelos::trace::source::TraceSourceEvent;
use zelos::trace::time::now_time_ns;
let measurement = source
.build_event("measurement")
.add_f64_field("value", None)
.build()?;
// Method 1: Current timestamp (automatic)
measurement.build().try_insert_f64("value", 1.0)?.emit()?;
// Method 2: Specific timestamp using emit_at
let specific_time_ns: i64 = 1_699_564_234_567_890_123;
measurement
.build()
.try_insert_f64("value", 2.0)?
.emit_at(specific_time_ns)?;
// Method 3: Calculated timestamp
let past_time_ns = now_time_ns() - (60 * 1_000_000_000); // 1 minute ago
measurement
.build()
.try_insert_f64("value", 3.0)?
.emit_at(past_time_ns)?;
// Method 4: Synchronized timestamps
async fn stream_synchronized(
sync_offset_ns: i64,
measurement: &TraceSourceEvent,
) -> Result<(), Box<dyn Error>> {
let mut interval = tokio::time::interval(Duration::from_millis(100));
loop {
interval.tick().await;
let timestamp_ns = now_time_ns() + sync_offset_ns;
let value = read_sensor(); // Your sensor reading function
measurement
.build()
.try_insert_f64("value", value)?
.emit_at_async(timestamp_ns)
.await?;
}
}
// Method 5: Replay historical data with its original timestamps
struct Record {
timestamp_ns: i64,
value: f64,
}
fn replay_historical_data(
records: &[Record],
measurement: &TraceSourceEvent,
) -> Result<(), Box<dyn Error>> {
for record in records {
measurement
.build()
.try_insert_f64("value", record.value)?
.emit_at(record.timestamp_ns)?;
}
Ok(())
}
measurement, err := source.BuildEvent("measurement").
AddFloat64Field("value", nil).
Build()
if err != nil {
log.Fatal(err)
}
// Method 1: Current timestamp (automatic)
builder, _ := measurement.Build().TryInsertFloat64("value", 1.0)
builder.Emit()
// Method 2: Specific timestamp using EmitAt
specificTimeNs := int64(1699564234567890123)
builder, _ = measurement.Build().TryInsertFloat64("value", 2.0)
builder.EmitAt(specificTimeNs)
// Method 3: Calculated timestamp
pastTimeNs := time.Now().Add(-1 * time.Minute).UnixNano()
builder, _ = measurement.Build().TryInsertFloat64("value", 3.0)
builder.EmitAt(pastTimeNs)
// Method 4: Synchronized timestamps
streamSynchronized := func(syncOffsetNs int64) {
ticker := time.NewTicker(100 * time.Millisecond)
defer ticker.Stop()
for range ticker.C {
timestampNs := time.Now().UnixNano() + syncOffsetNs
value := readSensor() // Your sensor reading function
b, _ := measurement.Build().TryInsertFloat64("value", value)
if err := b.EmitAt(timestampNs); err != nil {
log.Printf("Emit failed: %v", err)
}
}
}
// Method 5: Replay historical data with its original timestamps
type Record struct {
TimestampNs int64
Value float64
}
replayHistoricalData := func(records []Record) error {
for _, record := range records {
b, _ := measurement.Build().TryInsertFloat64("value", record.Value)
if err := b.EmitAt(record.TimestampNs); err != nil {
return err
}
}
return nil
}
Performance Tips¶
The tips below use Python. The same ideas apply in Rust and Go.
1. Define Schemas Upfront¶
A declared schema fixes field names, types, and units up front. Logging through the event handle also skips the per-call event lookup:
# Good: Schema defined once, reused many times
event = source.add_event("data", [
zelos_sdk.TraceEventFieldMetadata("value", zelos_sdk.DataType.Float64)
])
for i in range(1000000):
event.log(value=i * 0.1)
# Less efficient: looks the event up by name on every call
for i in range(1000000):
source.log("data", {"value": i * 0.1})
2. Use Appropriate Data Types¶
Choose the smallest data type that fits your needs:
# Use Int8 for small values (-128 to 127)
source.add_event("status", [
zelos_sdk.TraceEventFieldMetadata("state", zelos_sdk.DataType.Int8)
])
# Use Float32 when Float64 precision isn't needed
source.add_event("sensor", [
zelos_sdk.TraceEventFieldMetadata("temperature", zelos_sdk.DataType.Float32, "°C")
])
3. Group Related Fields¶
Log related fields in one event:
# Good: one event carries every field on one timestamp
imu.log(accel_x=ax, accel_y=ay, accel_z=az,
gyro_x=gx, gyro_y=gy, gyro_z=gz)
# Less efficient: two events, and the readings no longer share a timestamp
accel.log(x=ax, y=ay, z=az)
gyro.log(x=gx, y=gy, z=gz)
4. Log Blocks of Samples in One Call¶
When samples arrive in blocks, use log_many or log_batch instead of one log call per sample.
5. Control Timing Precisely¶
For high-frequency streaming, sleep until the next period instead of a fixed sleep() per loop.
Common Patterns¶
State Machine Monitoring¶
state_event = source.add_event("state", [
zelos_sdk.TraceEventFieldMetadata("current", zelos_sdk.DataType.UInt8),
zelos_sdk.TraceEventFieldMetadata("previous", zelos_sdk.DataType.UInt8),
zelos_sdk.TraceEventFieldMetadata("transition_time_ms", zelos_sdk.DataType.Float64)
])
# Add value table for readable state names
source.add_value_table("state", "current", {
0: "IDLE", 1: "INIT", 2: "RUNNING", 3: "ERROR"
})
source.add_value_table("state", "previous", {
0: "IDLE", 1: "INIT", 2: "RUNNING", 3: "ERROR"
})
For value tables in Rust and Go, see Enumerations.