Skip to content

Feature/binary datapoint ingest - #123

Merged
olavgg merged 5 commits into
mainfrom
feature/binary-datapoint-ingest
Sep 16, 2026
Merged

olavgg merged 5 commits into
mainfrom
feature/binary-datapoint-ingest

Conversation

@olavgg

@olavgg olavgg commented Sep 12, 2026 •

Copy link
Copy Markdown
Contributor

What this adds

A second datapoint ingest endpoint, POST /timeseries/data/binary, for the high-rate class. It
exists 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/data is 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:

JSON binary
API CPU 287.8 s 6.3 s 46x
Consumer CPU 64.2 s 6.4 s 10x
Bytes on the wire 5.00 GB 346 MB 14.4x
Ingest wall time 52.0 s 30.4 s 1.7x
API peak resident 3.8 GB 5.2 GB costs more

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:

Concurrent clients JSON binary
1 720,022/s 3,839,090/s
4 1,588,183/s 11,318,630/s
8 1,756,275/s 14,558,392/s

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,000
points 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 TEXT timeseries could not receive datapoints through the JSON endpoint at all. The value-type
switch had no TEXT arm, 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:

  • One value-type catalogue. The seven timeseries value types were written out in four places,
    two of which carried a comment asking the reader to keep the copies in sync. DatapointValueType
    is 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.
  • One latest-value cache. The compare-and-set that writes a series' latest value existed twice
    against the same Valkey key. Its rule now has direct tests, which it never had.
  • No series metadata cache. The binary path briefly cached what it knows about a series for
    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.
  • A valueType casing fix. The binary path reported the value type lower case to both
    WebSocket 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.
  • Request logging no longer stringifies a binary body; it records the word binary.
  • ClickHouse pinned to 26.8.2.7 in the test fixtures and the development stack.

Testing

  • Unit suites across every touched module.

    • datahub-e2e is new: integrationTest drives the platform through the Java SDK only and
      asserts 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:benchmark is the comparison above, re-runnable. datahub-e2e/README.md has the
      setup, the two product limits that have to be off for a run of that size, and the traps.
    • Integration tests against real Pulsar and ClickHouse for the endpoint, the consumer and the
      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; FrameLimits and ArrowSchemaCanon are the normative caps and
    schemas in code. SCALABILITY.md carries the decision and every measurement above, and the
    working notes that preceded it have been folded into it.

olavgg and others added 4 commits September 11, 2026 19:26
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>
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>
@olavgg
olavgg merged commit 5463d63 into main Sep 16, 2026
5 checks passed
@olavgg
olavgg deleted the feature/binary-datapoint-ingest branch September 16, 2026 17:54
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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants