Repository navigation
Feature/binary datapoint ingest - #123
Merged
Merged
Conversation
POST /timeseries/data/binary takes zstd-compressed Arrow IPC frames, one per value type, instead of JSON. FrameWriter builds a frame in the platform's contract: envelope, series directory, the canonical schema per value type, rows sorted and de-duplicated per series, values checked the way the JSON path checks them on the server. insert_datapoints_binary resolves series through /timeseries/byids into a per-service cache, cuts frames at the caps, compresses them in parallel and posts them, retrying a 429 or a 5xx and rebuilding once when the server reports a stale series. The blocking client mirrors it. The frame writer's tests run offline and read every frame back with arrow-rs; frames written here were also parsed by the platform's Java reader for all seven value types. test_datapoints_binary is the ignored live twin of test_datapoints over the same sixty days of values, and validate_daily_avg no longer keys its expectations on one external id. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Olav Gjerde <olav@intellistream.ai>
… clients The PyO3 bindings had no binary ingest at all, so Python could not reach the path the Rust core already implements. They now carry insert_datapoints_binary and insert_from_lists_binary on the sync service, with zstd_level defaulting to 9 as in the Rust and Java SDKs; insert_from_lists_binary takes the parallel timestamp and value sequences a DataFrame's index and column arrive in. python_tests/bench_json_vs_binary.py is the Python half of the comparison and reports the columns the platform's datahub-e2e benchmark does, so the clients can be read side by side. Measured at 2 million points over 10 float series: 282k points/s and 149 ms mean latency on JSON against 750k and 39 ms on binary. The gap is wider than Java's because more of Python's cost is per-call and in building the JSON, and the binary path moves that work into Rust. bench_json_vs_binary in the Rust test module is the third one, same columns again. Three fixes to the live test while making it pass against a real stack: the latest-value check compares within a few ULP, because one fixture value carries all 17 significant digits a f64 holds and came back one ULP off after crossing three systems as text; the cursor walk asks for limit 0, since an explicit limit caps the whole result and returns no cursor, which read as a series that stopped at its first 100k points; and the aggregates are compared against the JSON path rather than against validate_daily_avg's constants, which encode a zero-order-hold weighted average this platform does not compute and which the JSON twin fails on identically. Not done: the async bindings have no binary method yet, and the bindings' ValueType still covers only bigint, float and text. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Olav Gjerde <olav@intellistream.ai>
Chunking a large insert cloned each collection into the request body and then cloned the whole body again, so every datapoint was copied twice on its way out and a DatapointString is two heap Strings. Both clones were dead: neither value was read afterwards. Moving instead is worth about 6 percent on the JSON path. The benchmark now refuses to run without optimisation. Measured under a plain `cargo test` it reported this SDK as slower than the Java one on both paths, 339k against 720k points per second on JSON and 1.0M against 3.8M on binary, which is not a believable result and was entirely the dev profile's opt-level 0. In release the binary path reaches 8.1M points per second inside the call, 2.7 times Java. AGENTS.md says so too, next to the same warning for the PyO3 module. The Python benchmark gained an in-call rate alongside the wall clock. Its wall clock is dominated by generating points in interpreted code, so that number said more about the loop in the script than about the transport; in-call it is 344k on JSON against 1.31M on binary. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Olav Gjerde <olav@intellistream.ai>
A binary ingest call of up to 3.2 million points packs into one request and this changes nothing, but a larger one becomes several and the api validates each independently, so serialising them only added latency. At 40 million points in calls of 10 million: 3.81M points/s against 3.63M, 9.40M against 8.34M inside the call, and a p99 9 percent lower. Bounded by request_concurrency, default 4, so a very large call cannot open an unbounded number of connections. The benchmark can now select which paths to run, because a sweep over request sizes does not need the JSON one and it was the expensive part, and it reports an in-call rate beside the wall clock. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Olav Gjerde <olav@intellistream.ai>
JosteinGj
approved these changes
Sep 15, 2026
Conflict in src/timeseries/test.rs: validate_daily_avg takes main's source-derived expected averages over the branch's flattened constants. Main dropped truncate_10 with its last caller, but the branch's daily_aggregates still uses it, so it is restored beside that helper. test_datapoints_binary now also runs validate_daily_avg, since the constants its comment gave as the reason to skip it are gone. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This was referenced Sep 17, 2026
JosteinGj
added a commit
that referenced
this pull request
Sep 17, 2026
#123 was branched before #128 gave ResponseError a content_type field and merged after it, so main stopped compiling at the three places the binary path builds one. All three are raised by the SDK with no server response behind them (a 204 that fails to deserialize, a value refused locally, a series /byids cannot find), so each carries None. Signed-off-by: jgjesdal <jostein@intellistream.ai> Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What this adds
A second datapoint ingest endpoint,
POST /timeseries/data/binary, for the high-rate class. Itexists to move the platform's first scaling limit, the API parsing every value, off the API and
onto the client.
The client resolves each series once, checks and sorts the values, writes them as Arrow IPC
frames in the exact schema the ClickHouse table wants, compresses each frame with zstd, and posts
them. The API validates a frame and forwards its bytes unchanged. Frames travel compressed
through a Pulsar topic of their own and the consumer merges them into large ClickHouse blocks.
The JSON path at
POST /timeseries/datais untouched and remains the default.Arrow IPC rather than ClickHouse Native: the two parse at the same speed on these tables once a
frame holds a thousand points or more, so the choice went to the format that also suits the read
path, that a DataFrame round-trips through, and that any Arrow-capable tool can produce. The
frame envelope keeps a codec byte so a Native payload could be added later without a new media
type.
What it measured
100 million float32 points down each path, on a single 32-core development machine running the
services, Pulsar and ClickHouse together:
The API CPU is the result: accepting a hundred million points costs the JSON path 288 seconds of
a core and the binary path 6. Throughput gains less because the client is then the limit.
Scaling the clients out shows the difference more plainly. The JSON path saturates near 1.75
million points a second and adding clients does not move it, because they all queue behind the
same parsing. The binary path keeps scaling until the host runs out of cores:
Four clients is the balanced point: a billion points went through at 11.3 million a second with
the Pulsar backlog never passing a hundred messages, so ClickHouse absorbed the rows as fast as
the clients produced them. Eight and twelve clients reach higher rates but the backlog grows,
which means they are using Pulsar's burst absorption rather than sustaining.
Two honest caveats, both recorded in
SCALABILITY.md: the benchmark sends 10,000 to 40,000points per collection, which is the favourable shape for the JSON path, and the generated signal
is a slow sine plus noise with one signal per series. Giving every series identical values let
zstd compress across them and reported 0.52 bytes per point against the honest 3.46.
A defect this found
A
TEXTtimeseries could not receive datapoints through the JSON endpoint at all. The value-typeswitch had no
TEXTarm, so every text datapoint fell through to the default and threw"Unsupported value type: TEXT", which the caller saw as a 500. Everything around it already
assumed text worked: the per-collection cap checks it, the quota counts it, and the consumer
stores it. Fixed here, with a unit test.
It was found by the new end-to-end test, which sends the same points down both paths and compares
what ClickHouse stored.
Also in this branch
Consolidations that the new code made visible, each with tests:
two of which carried a comment asking the reader to keep the copies in sync.
DatapointValueTypeis now the single definition; a parity test replays the value-type inserts of every migration
and compares them with the enum, so an eighth type added to either side fails the build.
against the same Valkey key. Its rule now has direct tests, which it never had.
fifteen seconds. The dataset id is an input to the write ACL, so a stale entry authorises
against the dataset a series has been moved out of, and its invalidation methods were never
called. Replaced with a projection query that reads four columns instead of hydrating entities.
valueTypecasing fix. The binary path reported the value type lower case to bothWebSocket paths where the JSON path reports it upper case, so a client could not treat them
alike. The two tests covering it had locked in the wrong spelling.
binary.Testing
Unit suites across every touched module.
datahub-e2eis new:integrationTestdrives the platform through the Java SDK only andasserts that the same points sent as JSON and as frames land identically. A bulk series is
compared by an ordered checksum computed inside ClickHouse rather than by row count, so a
value mangled by one path's encoding cannot pass, and all seven value types are then sent
both ways and their stored values compared against each other. The assertion is deliberately
that a reader cannot tell which path wrote the row, not that storage echoes the input back.
datahub-e2e:benchmarkis the comparison above, re-runnable.datahub-e2e/README.mdhas thesetup, the two product limits that have to be off for a run of that size, and the traps.
Arrow round trip.
Documentation
The wire contract for third-party producers, the SDK methods and the new limits are documented
separately in the docs sites;
FrameLimitsandArrowSchemaCanonare the normative caps andschemas in code.
SCALABILITY.mdcarries the decision and every measurement above, and theworking notes that preceded it have been folded into it.