A Tiny SQL Engine for JSON Streams Fits Between jq and a Stream Processing Cluster
Your service writes one JSON object per line to stdout, and you want the average cpu for host web-3 over the last five minutes. jq can’t answer that, because it forgets each object once it has printed it and has no idea what “five minutes” means. You can fake both with jq -n 'reduce inputs as $e (...)' plus a timestamp comparison, and you’ll have a script nobody wants to touch twice. DuckDB answers the question in one line, but only against a file as it exists right now. ksqlDB, Materialize, RisingWave, Arroyo and Flink SQL answer it continuously, once you’ve deployed a service and agreed to operate it.
The gap between those is a library with a CLI on top. It takes JSON objects as they arrive, works out the columns itself, and answers windowed SQL without keeping every event. Picture tail -f app.json | tailq "SELECT ..." on a laptop, inside a test suite, on an edge box, or in front of an agent that emits telemetry. (tailq is a placeholder for a tool that doesn’t exist.) The trick that makes it small is also its main limit: it keeps aggregates instead of events, so it can only answer the questions you registered before the data arrived, plus whatever fits in a bounded buffer of recent raw events.
Neighbors cover pieces of this. lnav runs SQL over log files you already have. sqlite-utils loads JSON lines into SQLite, and from there json_extract plus a GROUP BY answers most one-off questions, at the cost of storing every row. Fluent Bit has a SQL stream processor, but it lives inside a log shipper, so you adopt the shipper to get it. Try those first. They stop helping when you want a window that stays current while memory stays flat, or when you want the whole thing inside your own process.
Keep the Answers, Not the Events
Ingest has one job per line: parse it, pull out the paths the registered queries touch, and update some state. Nothing else gets stored. For each combination of query, window and group, the state is a handful of numbers: a count, a sum, a minimum and a maximum. The mean comes from the sum and the count, so it merges cleanly, which pays off once you add hopping windows: keep one partial aggregate per pane (a slice of timeline the size of the hop) and build each window by merging its panes. Each event then updates one pane instead of every window that overlaps it.
Registered queries are the cheap path. For everything else there’s a ring buffer holding the last N events or megabytes. A SELECT with no WINDOW clause runs over that buffer and nothing older. That’s the contract, and the tool should print it with every ad hoc answer: “covers 09:02:11 to 09:07:48, 50,000 events”. An engine that quietly answers over a fraction of the data is worse than one that refuses.
$ tail -f app.jsonl | tailq --time ts --lateness 10s \
"SELECT host, avg(cpu) AS cpu, max(cpu) AS peak, count(*) AS n
FROM stream
WHERE env = 'prod'
GROUP BY host
WINDOW TUMBLING 5m"
window_start host cpu peak n
2026-10-05T09:00:00Z web-1 41.2 88.0 882
2026-10-05T09:00:00Z web-3 77.9 99.1 864
tailq: 2 events arrived after window 09:00 closed (lateness 10s), dropped
tailq: cpu: 14 string values, 9 coerced to numbers, 5 rejected (first: "n/a")
2026-10-05T09:05:00Z web-1 39.8 71.5 871
2026-10-05T09:05:00Z web-3 80.3 98.7 860
A window prints once, when it closes, so the first row lands at most five minutes plus the lateness allowance after the first event. The library form is the same thing without the pipe:
s = Stream(ring_events=50_000, lateness="10s", spill="metrics.db")
q = s.continuous("SELECT host, avg(cpu) FROM stream GROUP BY host WINDOW TUMBLING 5m")
s.write({"ts": "2026-10-05T09:00:03Z", "host": "web-3", "cpu": 77.9})
q.latest() # rows from the most recent closed window
s.query("SELECT count(*) FROM stream WHERE host = 'web-3'") # ring buffer only
Keep the SQL subset small enough to parse by hand: a select list of columns, literals, arithmetic and aggregate calls, then WHERE, GROUP BY and a WINDOW clause. No joins, subqueries or CTEs. WINDOW 5m isn’t standard SQL, so this is a dialect, and the docs should say so on page one. Resist the shortcut of handing the ring buffer to SQLite for ad hoc queries. You’d end up with two evaluators, and the same expression would give different answers depending on which one ran it (SQLite divides integers as integers, for one). SQLite belongs on the spill side, as storage.
Type Drift Needs a Policy
Real streams drift. cpu arrives as 12.5 for a week, then a new agent version starts sending "12.5%", then somebody’s script sends "n/a". A schemaless engine has to pick a behavior, and the worst pick is the quiet one, because an average over the values that happened to parse still looks plausible.
The policy worth shipping: each path keeps type counts as events arrive (numbers, strings, booleans, nulls, objects, arrays). A query that uses cpu as a number gets numeric values as they are, numeric-looking strings coerced (a trailing percent sign gets stripped, nothing fancier), and everything else treated as NULL. Every coercion and every rejection is counted, and the counts print next to the answer. A --schema flag shows the type mix per path, which is the first thing you’ll want when a dashboard number looks off.
Time Is the Hard Part
Event time is when something happened. Arrival time is when you heard about it. Windows should follow event time, which means deciding when a window is finished even though a late event could still turn up. The standard answer is a watermark: the largest event time seen so far, minus an allowed lateness. A window closes when the watermark passes its end. Events older than the watermark get counted and dropped. Simple, with two sharp edges.
The first is a bad clock. One device reporting a timestamp three days ahead drags the watermark forward, and every honest event after that looks late. Clamp events that sit far beyond the arrival time, keep them from advancing the watermark, and count them separately. The second is a quiet stream. No events means no watermark movement, so the last window of the evening stays open until morning. Fall back to the wall clock after an idle timeout, and say so in the output when that fires.
High-cardinality GROUP BY is the memory risk. Group by user id on a busy stream and the state grows with the number of users, per window. Cap the groups per query window, fold the overflow into one __other__ group, and print how many keys landed there. Space-Saving is the usual answer when you’d rather keep the heaviest keys than lump them together. Exact distinct counts and percentiles can’t survive at this scale, so they arrive later as sketches (HyperLogLog for distinct counts, KLL or a t-digest for percentiles), which is a design problem of its own. Sketches cost kilobytes per group per window, so they stay opt-in per query.
Persistence needs an honest answer too. Open windows live in RAM and die with the process. Closed windows go to one SQLite file, one row per query, window and group, in a transaction that also records the input offset to resume from (the earliest event any open window still depends on). After a crash with file input, the tool re-reads from that offset and rebuilds the open windows exactly. A pipe has no offset, so the windows that were open at the crash are partial, and the output should flag them. Closed windows are never lost, and nothing here claims exactly-once. It’s the same forgetting-on-purpose trade as a telemetry database that keeps rollups instead of raw rows, moved to the front door. Precomputing attacks it from the storage side: a policy compiles to SQLite triggers so every insert keeps the declared answers current.
Version 0.1 Does Less Than You Want
The first version is a library plus a CLI, with JSON Lines in, a ring buffer, tumbling windows, five aggregates (count, sum, min, max, avg) and spill to a SQLite file. Hopping windows, sketches and a persistent ring buffer wait. It refuses joins between streams, exactly-once delivery and clustering, because each one turns the tool into the stream processor it was meant to sit beside.
Parsing will be the cost that shows up in a profile, while the aggregation barely registers. Every line gets parsed once, and the engine only needs the paths its queries mention, so a parser that skips untouched fields (the idea behind simdjson’s On-Demand API) matters more than clever state. If the same big documents get read over and over instead, persisting them as indexed binary is the sibling problem. If you want the raw lines kept and queryable later, a single-file log database is the tool for that, and this one shouldn’t grow into it.
Start with tumbling windows and five aggregates. The rest can wait.