Skip to content

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_many rejects 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_many skips a key that is not in the schema instead of raising, so a typo loses that value without an error.
  • Use flat dicts. log_many does not log nested dicts.
  • None is not a valid field value here. log_many raises ValueError for 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_batch reads the first column as time whatever its name, and does not check its type. A batch with another time type, such as int64 or timestamp("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, and timestamp("ns").
  • log_batch drops columns of any other Arrow type without an error. pandas text columns become large_string, so cast them to string first.
  • The first batch for an event name registers its schema, without units. To give fields units, register the event with add_event first. log_batch then 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_batch sorts 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")
])

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.

Sensor Array Streaming

# Define schema for array of sensors
sensor_array = source.add_event("array", [
    zelos_sdk.TraceEventFieldMetadata(f"sensor_{i}", zelos_sdk.DataType.Float32)
    for i in range(16)
])

# Log all sensors with single timestamp
sensor_array.log(**{f"sensor_{i}": values[i] for i in range(16)})