Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,17 @@ cargo test <name> # substring match on test name
cargo test -- --ignored # run tests marked #[ignore] (e.g. long-running datapoint tests)
cargo test <path>::tests::<name> # e.g. `events::tests::test_events_full`
cargo test -- --nocapture # show println! from tests (the SDK prints response bodies)
cargo test --release <bench name> # ALWAYS --release for anything timed (see below)
./run_python_tests.sh # Python-bindings suite (rebuilds the PyO3 module first — see below)
```

**Never time anything under a plain `cargo test`.** It builds with the `dev` profile at
`opt-level = 0`. The `bench_json_vs_binary` comparison first reported Rust as *slower than the
Java SDK* on both ingest paths that way, 339k against 720k points per second on JSON and 1.0M
against 3.8M on binary; the whole result was the missing `--release`. That test now panics
rather than run unoptimised, and the same applies to the PyO3 module, which needs
`maturin develop --release` before any Python timing means anything.

Most tests are integration tests that call a live backend via `create_api_service()`. They read configuration from a local `.env` file (gitignored). Required:

- `BASE_URL` — backend root, e.g. `http://localhost:8081`
Expand Down Expand Up @@ -115,6 +123,10 @@ Synchronous mirror of the async API behind the `blocking` cargo feature — the

When a datapoint/event send can't get through, ingestion spools to a segmented, zstd-compressed NDJSON log on disk and flushes automatically on a later ingest call. Invariants to preserve: memory use is bounded by a single segment (plain append-only active segment, zstd-sealed at ~50 MiB rollover via temp file + atomic rename, drained oldest-first one segment at a time); bounded by time retention (whole segments past the window dropped, expired records skipped on read) and a size cap (oldest segment deleted); a torn trailing line from an unclean shutdown is skipped on read. Each on-disk line is `<epoch_millis>\t<json>`; the spool is content-agnostic.

### Binary datapoint ingest (`src/timeseries/binary.rs`)

`TimeSeriesService::insert_datapoints_binary` is the second ingest path, to `POST /timeseries/data/binary`: the same `DatapointsCollection<DatapointString>` input, resolved through `/timeseries/byids` (cached per service instance; needs read access on the dataset), checked against each series' value type locally, cut into Arrow IPC frames at the contract's caps, zstd-compressed per frame (mandatory: level 1, 3 or 9, default 9) and posted through `execute_post_bytes_request`. `FrameWriter` builds one frame and is public; `cut_into_writers` and `pack_requests` hold the caps. The byte layout is the platform's `binary_datapoints_format.md`. The arrow-rs crates (`arrow-array`, `arrow-schema`, `arrow-ipc`) exist for this path and are the seed of the Arrow read path. Not in the Python bindings yet. The ignored `timeseries::tests::test_datapoints_binary` is the live twin of `test_datapoints` and needs a backend that serves the endpoint; the writer's own tests in `binary.rs` run offline.

### The `ApiServiceProvider` trait (`src/generic.rs`)

Every subservice implements `ApiServiceProvider`, which owns the HTTP plumbing: token acquisition, `execute_get_request`, `execute_post_request`, `execute_file_upload_request`, `execute_get_stream_request`. Subservice methods should go through these helpers rather than calling `reqwest` directly — a few early methods (e.g. `TimeSeriesService::list`) still bypass the trait and should be migrated when touched.
Expand Down
3 changes: 3 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,9 @@ geojson = "1"
# Reading (never verifying) the payload of a JWT the SDK already holds, to explain an
# otherwise-unexplained 401 — see `auth_diagnostics`.
base64 = "0.22"
arrow-array = "59.3"
arrow-schema = "59.3"
arrow-ipc = "59.3"


#[lib]
Expand Down
16 changes: 16 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,22 @@ sealed at a ~50 MiB rollover and drained one segment at a time — so even a mul
spool never loads into memory, and a torn trailing line from an unclean shutdown is skipped
on read.

## Binary datapoint ingest

`time_series.insert_datapoints_binary(&collections, &BinaryIngestOptions::default())` takes
the same collections as `insert_datapoints` and sends them through
`POST /timeseries/data/binary`: each series is resolved once to its id and value type and
cached, values are checked against that type locally, sorted and de-duplicated, cut into Arrow
IPC frames at the contract's caps (100 000 points per numeric frame, 10 000 per text or mixed
frame, 32 frames per request), compressed with zstd (level 9 by default, 1 and 3 are the other
choices) and posted. A 204 means every frame was accepted.

A series that does not exist is a 404 before anything is sent, a value that does not fit its
type is a 422, and a 429 or a 5xx is retried. Resolving by external id goes through
`/timeseries/byids`, so the caller needs read access on the dataset as well as write access.
The durable spool does not cover this path. `binary::FrameWriter` is public for producers that
build frames themselves.

## Python bindings

`datahub_python_bindings/` wraps this SDK as the Python package `intellistream-datahub-sdk` (import name
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -635,6 +635,18 @@ class TimeSeriesServiceSync:
"""

def insert_datapoints(self, input: list[DatapointsCollectionString]) -> list[str]: ...
def insert_datapoints_binary(
self,
input: list[DatapointsCollectionString],
zstd_level: int | None = None,
) -> list[str]: ...
def insert_from_lists_binary(
self,
timestamps: list[datetime.datetime],
values: list[float],
ts: Identifiable,
zstd_level: int | None = None,
) -> list[str]: ...
def insert_from_lists(
self,
timestamps: list[datetime.datetime],
Expand Down
49 changes: 49 additions & 0 deletions datahub_python_bindings/src/timeseries/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -469,6 +469,55 @@ pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> {
Ok(())
}

/// Options for the binary ingest bindings. zstd level 1, 3 or 9; 9 is the default, as in the
/// Rust and Java SDKs, because the client is what pays for it.
pub(crate) fn binary_options(
zstd_level: Option<i32>,
) -> intellistream_datahub_sdk::timeseries::BinaryIngestOptions {
use intellistream_datahub_sdk::timeseries::BinaryIngestOptions;
match zstd_level {
Some(level) => BinaryIngestOptions::new().zstd_level(level),
None => BinaryIngestOptions::default(),
}
}

/// Parallel timestamp and value sequences for one series, which is the shape a DataFrame's
/// index and one of its columns arrive in, turned into the wire collection.
pub(crate) fn lists_to_collection(
timestamps: Vec<Bound<'_, PyAny>>,
values: Vec<f64>,
ts: Identifiable,
) -> PyResult<DatapointsCollection<DatapointString>> {
if timestamps.len() != values.len() {
return Err(PyValueError::new_err(format!(
"timestamps and values must be the same length, got {} and {}",
timestamps.len(),
values.len()
)));
}
let datapoints: Vec<DatapointString> = timestamps
.into_iter()
.zip(values)
.map(|(timestamp, value)| {
Ok(DatapointString {
timestamp: crate::datetime::py_datetime_to_utc(&timestamp)?
.timestamp_millis()
.to_string(),
value: value.to_string(),
})
})
.collect::<PyResult<Vec<_>>>()?;
let id_collection = ts.id_collection();
Ok(DatapointsCollection {
datapoints,
next_cursor: None,
id: id_collection.id,
external_id: id_collection.external_id,
unit: None,
unit_external_id: None,
})
}

/// Reverse lookup: the events that reference this timeseries. (The other direction —
/// `Event.related_resource_nodes()` — resolves an event's resources.) Available only on
/// timeseriess returned by the API; calling on a locally-constructed one raises.
Expand Down
47 changes: 47 additions & 0 deletions datahub_python_bindings/src/timeseries/sync_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,53 @@ impl PyTimeSeriesServiceSync {
Ok(result.get_items().clone())
})
}
/// `POST /timeseries/data/binary`: the same collections as `insert_datapoints`, sent as
/// zstd-compressed Arrow frames. `zstd_level` is 1, 3 or 9 and defaults to 9.
#[pyo3(signature = (input, zstd_level=None))]
fn insert_datapoints_binary<'py>(
&self,
py: Python<'py>,
input: Vec<PyDatapointsCollectionString>,
zstd_level: Option<i32>,
) -> PyResult<Vec<String>> {
let service = self.api_service.clone();
let vec: Vec<DatapointsCollection<DatapointString>> =
input.into_iter().map(|item| item.into()).collect();
let wrapper = DataWrapper::<DatapointsCollection<DatapointString>>::from_vec(vec);
let options = binary_options(zstd_level);
py.detach(|| {
let result = self
.runtime
.block_on(service.time_series.insert_datapoints_binary(&wrapper, &options))
.map_err(|e| crate::datahub_err(e))?;
Ok(result.get_items().clone())
})
}

/// The binary twin of `insert_from_lists`: parallel timestamp and value sequences for one
/// series, which is the shape a DataFrame column pair arrives in.
#[pyo3(signature = (timestamps, values, ts, zstd_level=None))]
fn insert_from_lists_binary<'py>(
&self,
py: Python<'py>,
timestamps: Vec<Bound<'py, PyAny>>,
values: Vec<f64>,
ts: Identifiable,
zstd_level: Option<i32>,
) -> PyResult<Vec<String>> {
let service = self.api_service.clone();
let collection = lists_to_collection(timestamps, values, ts)?;
let wrapper = DataWrapper::<DatapointsCollection<DatapointString>>::from_vec(vec![collection]);
let options = binary_options(zstd_level);
py.detach(|| {
let result = self
.runtime
.block_on(service.time_series.insert_datapoints_binary(&wrapper, &options))
.map_err(|e| crate::datahub_err(e))?;
Ok(result.get_items().clone())
})
}

fn insert_from_lists<'py>(
&self,
py: Python<'py>,
Expand Down
166 changes: 166 additions & 0 deletions python_tests/bench_json_vs_binary.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,166 @@
# JSON against binary datapoint ingest, through the Python SDK.
#
# Not a pytest: it needs a live platform and takes minutes, so it is a script you run when you
# want the numbers. The Java equivalent lives in the platform repo's datahub-e2e module and
# reports the same columns, so the two are comparable.
#
# python python_tests/bench_json_vs_binary.py --points 10000000
#
# Reads the same .env the SDK does. The api must have its daily quota and rate limiter off, or
# it will refuse a run of this size:
# -Ddatahub.limits.quota.enabled=false -Ddatahub.limits.rate.enabled=false
"""Measure the two ingest paths from Python and print a comparison."""

from __future__ import annotations

import argparse
import datetime as dt
import math
import os
import resource
import time
import urllib.request

from intellistream_datahub_sdk import DataHubClient, TimeSeries

START = dt.datetime(2025, 1, 1, tzinfo=dt.timezone.utc)


def generate(count: int, offset: int, series_index: int) -> tuple[list[dt.datetime], list[float]]:
"""A slow sine plus noise, one signal per series.

Identical series would let zstd compress the repetition across them and report a wire size
no real fleet of sensors would produce.
"""
base = 150.0 + series_index * 0.7
phase = series_index * 0.37
timestamps = []
values = []
for i in range(count):
t = offset + i
timestamps.append(START + dt.timedelta(seconds=t))
# Deterministic, so two runs generate the same bytes.
z = (t + series_index * 6364136223846793005) * 6364136223846793005 & 0xFFFFFFFFFFFFFFFF
z ^= z >> 33
noise = ((z >> 40) / float(1 << 24)) - 0.5
values.append(base + 20.0 * math.sin(t / 600.0 + phase) + noise)
return timestamps, values


def clickhouse_count(ids: list[int]) -> int:
url = (
os.environ.get("CLICKHOUSE_URL", "http://localhost:18123")
+ "/?user=" + os.environ.get("CLICKHOUSE_USER", "foobar")
+ "&password=" + os.environ.get("CLICKHOUSE_PASSWORD", "changeme")
)
db = os.environ.get("CLICKHOUSE_DB", "foo")
sql = (
f"SELECT count() FROM {db}.datapoints_float WHERE timeseries_id IN "
f"({','.join(str(i) for i in ids)})"
)
with urllib.request.urlopen(url, sql.encode()) as response:
return int(response.read().decode().strip())


def peak_rss_mb() -> float:
# ru_maxrss is kilobytes on Linux.
return resource.getrusage(resource.RUSAGE_SELF).ru_maxrss / 1024.0


def run(client, label: str, binary: bool, points: int, series: int, chunk: int) -> dict:
run_id = f"pybench_{int(time.time())}_{'bin' if binary else 'json'}"
external_ids = [f"{run_id}_{i}" for i in range(series)]
internal_ids = []
for external_id in external_ids:
created = client.timeseries.create(
[TimeSeries(name=external_id, external_id=external_id, value_type="float", unit="celsius")]
)
internal_ids.append(created[0].id)

per_series = chunk // series
latencies = []
sent = 0
offset = 0
started = time.perf_counter()
while sent < points:
this_chunk = min(chunk, points - sent)
per_series = max(1, this_chunk // series)
for index, external_id in enumerate(external_ids):
timestamps, values = generate(per_series, offset, index)
call = time.perf_counter()
if binary:
client.timeseries.insert_from_lists_binary(timestamps, values, external_id)
else:
client.timeseries.insert_from_lists(timestamps, values, external_id)
latencies.append(time.perf_counter() - call)
sent += per_series * series
offset += per_series
elapsed = time.perf_counter() - started
print(f" {label}: {sent:,} / {points:,} points, {sent / elapsed:,.0f} pts/s", flush=True)
ingest_seconds = time.perf_counter() - started

settle_start = time.perf_counter()
deadline = settle_start + 1800
while time.perf_counter() < deadline:
if clickhouse_count(internal_ids) >= sent:
break
time.sleep(1)
else:
raise SystemExit(f"{label}: only {clickhouse_count(internal_ids)} of {sent} rows became readable")
settle_seconds = time.perf_counter() - settle_start

client.timeseries.delete(external_ids)
latencies.sort()
return {
"path": label,
"points": sent,
"ingest_seconds": ingest_seconds,
"settle_seconds": settle_seconds,
"points_per_second": sent / ingest_seconds,
# Wall clock includes generating the points in interpreted Python, which dominates it.
# This is the same points divided by the time actually spent inside the SDK calls, so
# it says what the transport did rather than what the loop above did.
"points_per_second_in_call": sent / sum(latencies),
"requests": len(latencies),
"latency_mean_ms": sum(latencies) / len(latencies) * 1000,
"latency_p50_ms": latencies[len(latencies) // 2] * 1000,
"latency_p99_ms": latencies[min(len(latencies) - 1, int(len(latencies) * 0.99))] * 1000,
"client_peak_rss_mb": peak_rss_mb(),
}


def main() -> None:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--points", type=int, default=10_000_000)
parser.add_argument("--series", type=int, default=100)
parser.add_argument("--chunk", type=int, default=1_000_000)
args = parser.parse_args()

env_file = os.environ.get("DATAHUB_ENV_FILE", ".env")
client = DataHubClient.from_envfile(env_file)
print(f"\n=== {args.points:,} points across {args.series:,} series, float ===", flush=True)
results = [
run(client, "JSON", False, args.points, args.series, args.chunk),
run(client, "binary", True, args.points, args.series, args.chunk),
]

rows = [
("ingest wall time (s)", "{:.1f}", "ingest_seconds"),
("points per second", "{:,.0f}", "points_per_second"),
("points per second in-call", "{:,.0f}", "points_per_second_in_call"),
("settle to readable (s)", "{:.1f}", "settle_seconds"),
("requests", "{:,.0f}", "requests"),
("latency mean (ms)", "{:,.0f}", "latency_mean_ms"),
("latency p50 (ms)", "{:,.0f}", "latency_p50_ms"),
("latency p99 (ms)", "{:,.0f}", "latency_p99_ms"),
("client peak RSS (MB)", "{:,.0f}", "client_peak_rss_mb"),
]
print(f"\n=== datapoint ingest from Python: JSON against binary ===")
print(f"{args.points:,} points across {args.series:,} series, float\n")
print("{:<26}{:>18}{:>18}".format("metric", *[r["path"] for r in results]))
for label, fmt, key in rows:
print("{:<26}{:>18}{:>18}".format(label, *[fmt.format(r[key]) for r in results]))


if __name__ == "__main__":
main()
8 changes: 7 additions & 1 deletion src/blocking.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ use crate::relations::{EdgeProxy, RelForm, RelTypeForm, RelationshipType};
use crate::resources::{
RelatedResourcesForm, Resource, ResourceFilter, ResourceNetwork, ResourceUpdate,
};
use crate::timeseries::{TimeSeries, TimeSeriesFilter, TimeSeriesUpdateCollection};
use crate::timeseries::{BinaryIngestOptions, TimeSeries, TimeSeriesFilter, TimeSeriesUpdateCollection};
use crate::unit::Unit;

/// Generate blocking methods that delegate to the same-named async method on one of
Expand Down Expand Up @@ -168,6 +168,7 @@ impl TimeSeriesService {
fn search_by_query(query: &str) -> Result<DataWrapper<TimeSeries>, ResponseError>;
fn insert_datapoint(id: Option<u64>, external_id: Option<String>, timestamp: DateTime<Utc>, value: String) -> Result<DataWrapper<String>, ResponseError>;
fn insert_datapoints(json: &mut DataWrapper<DatapointsCollection<DatapointString>>) -> Result<DataWrapper<String>, ResponseError>;
fn insert_datapoints_binary(json: &DataWrapper<DatapointsCollection<DatapointString>>, options: &BinaryIngestOptions) -> Result<DataWrapper<String>, ResponseError>;
fn retrieve_datapoints(json: &DataWrapper<RetrieveFilter>) -> Result<DataWrapper<DatapointsCollection<Datapoint>>, ResponseError>;
fn delete_datapoints(json: &DataWrapper<DeleteFilter>) -> Result<DataWrapper<String>, ResponseError>;
fn retrieve_latest_datapoint(json: &DataWrapper<IdAndExtId>) -> Result<DataWrapper<DatapointsCollection<Datapoint>>, ResponseError>;
Expand All @@ -177,6 +178,11 @@ impl TimeSeriesService {
pub fn buffered_count(&self) -> u64 {
self.api.time_series.buffered_count()
}

/// Already synchronous on the async service; passed through directly.
pub fn evict_binary_series_cache(&self) {
self.api.time_series.evict_binary_series_cache()
}
}

/// Blocking counterpart of [`crate::ResourceService`].
Expand Down
Loading
Loading