<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: Ankur Kumar Pandey</title>
    <description>The latest articles on DEV Community by Ankur Kumar Pandey (@ankurpaan).</description>
    <link>https://dev.to/ankurpaan</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F4118447%2Fc73d8bfb-562a-45ae-b7e9-c2708aea5e2d.jpg</url>
      <title>DEV Community: Ankur Kumar Pandey</title>
      <link>https://dev.to/ankurpaan</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/ankurpaan"/>
    <language>en</language>
    <item>
      <title>rekuiper 0.500: Moving the Engine Hot Path to RAM and Finding its Exact Physical Limits (up to 200k msg/s)</title>
      <dc:creator>Ankur Kumar Pandey</dc:creator>
      <pubDate>Mon, 21 Sep 2026 23:39:35 +0000</pubDate>
      <link>https://dev.to/ankurpaan/rekuiper-0500-moving-the-engine-hot-path-to-ram-and-finding-its-exact-physical-limits-up-to-200k-53di</link>
      <guid>https://dev.to/ankurpaan/rekuiper-0500-moving-the-engine-hot-path-to-ram-and-finding-its-exact-physical-limits-up-to-200k-53di</guid>
      <description>&lt;p&gt;When we published our earlier benchmarks for &lt;strong&gt;&lt;a href="https://github.com/ankur-paan/rekuiper" rel="noopener noreferrer"&gt;rekuiper&lt;/a&gt;&lt;/strong&gt; (our Rust reimplementation of LF Edge eKuiper for edge gateways and IoT hubs), the numbers answered our first question: on five real-world MQTT workloads, memory stayed bounded between &lt;strong&gt;5 and 10 MB&lt;/strong&gt; while Go-based engines climbed to hundreds of megabytes or failed.&lt;/p&gt;

&lt;p&gt;In that test, rekuiper sustained &lt;strong&gt;100,000 messages per second on a single pinned core&lt;/strong&gt;. At offered 200,000 msg/s, however, it dropped packets or backlogged upstream. &lt;/p&gt;

&lt;p&gt;Saying "100k passed and 200k failed" left a massive 100,000 msg/s blind spot. Is the true ceiling 105k? 140k? 195k? And &lt;em&gt;why&lt;/em&gt; did it choke at higher rates when the CPU wasn't fully maxed out on every workload?&lt;/p&gt;

&lt;p&gt;When we profiled the bottleneck under burst loads, the culprit wasn't stream parsing or window math. &lt;strong&gt;It was the disk.&lt;/strong&gt;&lt;/p&gt;

&lt;p&gt;Even when an edge stream engine's data plane runs in memory, typical implementations still touch disk for metadata: checking SQLite tables for stream and rule definitions, resolving auth keys, reading configuration YAMLs, or persisting rule state changes. On a developer NVMe drive, SQLite lookups take microseconds. On an industrial gateway running off slow eMMC flash or a microSD card, a batch of flash writes causes I/O wait spikes that stall the Tokio runtime thread.&lt;/p&gt;

&lt;p&gt;For &lt;strong&gt;v0.500-beta&lt;/strong&gt;, we rebuilt the engine's internal catalog and hot paths around an in-memory Redis-style architecture, tuned our storage engine to be zero-contention, and then ran a hierarchical search with 1,000 msg/s resolution under a strictly bounded Mosquitto broker to find the exact physical limit of every workload.&lt;/p&gt;

&lt;p&gt;Here is what we built, what we measured, and where single-core stream processing actually hits physical walls.&lt;/p&gt;




&lt;h2&gt;
  
  
  1. The Bottleneck: Why Disk Kills Edge Ingest
&lt;/h2&gt;

&lt;p&gt;In eKuiper and similar edge software, metadata (streams, rules, schemas, auth) is stored in a key-value store or an embedded SQLite database. &lt;/p&gt;

&lt;p&gt;Under moderate load (5k–20k msg/s), querying SQLite to validate a rule or look up a schema is unnoticeable. But when 100,000+ messages per second hit an edge machine:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;Lock Contention&lt;/strong&gt;: When multiple concurrent connections or rules access SQLite simultaneously, lock acquisition overhead adds tail latency.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Flash Storage Stalls&lt;/strong&gt;: Embedded devices (Raspberry Pi, industrial gateways) do not have enterprise SSD write queues. If the engine writes checkpoint state or logs to flash while simultaneously servicing network buffers, the OS I/O scheduler blocks threads.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Queue Backpressure&lt;/strong&gt;: When a stream processor stalls for even 10 milliseconds waiting on disk, a 150k msg/s MQTT ingress generates a 1,500-message backlog in the socket buffer. Under QoS 0, Mosquitto drops those packets.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;If we wanted to push edge ingestion beyond 100k msg/s on a single core, &lt;strong&gt;the hot path could not touch the filesystem at all.&lt;/strong&gt;&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2F44m2ypjeyl8a233rd46b.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2F44m2ypjeyl8a233rd46b.png" alt=" " width="800" height="498"&gt;&lt;/a&gt;&lt;/p&gt;




&lt;h2&gt;
  
  
  2. What Changed in v0.500
&lt;/h2&gt;

&lt;h3&gt;
  
  
  Redis-Style In-Memory Catalog (&lt;code&gt;MemoryCatalog&lt;/code&gt;)
&lt;/h3&gt;

&lt;p&gt;We separated metadata persistence from metadata reads:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight rust"&gt;&lt;code&gt;&lt;span class="c1"&gt;// crates/rekuiper-core/src/catalog.rs&lt;/span&gt;
&lt;span class="k"&gt;pub&lt;/span&gt; &lt;span class="k"&gt;struct&lt;/span&gt; &lt;span class="n"&gt;MemoryCatalog&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
    &lt;span class="n"&gt;streams&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nn"&gt;parking_lot&lt;/span&gt;&lt;span class="p"&gt;::&lt;/span&gt;&lt;span class="n"&gt;RwLock&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="n"&gt;HashMap&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nb"&gt;String&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;StreamDefinition&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&amp;gt;&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
    &lt;span class="n"&gt;tables&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nn"&gt;parking_lot&lt;/span&gt;&lt;span class="p"&gt;::&lt;/span&gt;&lt;span class="n"&gt;RwLock&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="n"&gt;HashMap&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nb"&gt;String&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;TableDefinition&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&amp;gt;&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
    &lt;span class="n"&gt;rules&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nn"&gt;parking_lot&lt;/span&gt;&lt;span class="p"&gt;::&lt;/span&gt;&lt;span class="n"&gt;RwLock&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="n"&gt;HashMap&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nb"&gt;String&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;RuleDefinition&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&amp;gt;&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
&lt;span class="p"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;ul&gt;
&lt;li&gt;On daemon startup, the catalog hydrates all streams, tables, and active rules from SQLite into memory once.&lt;/li&gt;
&lt;li&gt;During execution, the stream bus, rule evaluators, and REST query endpoints read directly from &lt;code&gt;RwLock&amp;lt;HashMap&amp;gt;&lt;/code&gt;, completing lookups in nanoseconds with zero system calls and zero disk I/O.&lt;/li&gt;
&lt;li&gt;Rule mutations update memory first and asynchronously commit to SQLite in the background.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  Zero-Disk Hot Path &amp;amp; Caching
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Auth Token Cache&lt;/strong&gt;: Public RSA keys used for JWT signature verification (&lt;code&gt;KUIPER_AUTH_PUBLIC_KEY_FILE&lt;/code&gt;) are parsed and cached in memory. Ingest requests no longer read public keys off disk per request.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Config &amp;amp; Schema Cache&lt;/strong&gt;: &lt;code&gt;/etc&lt;/code&gt; configuration overlays, source definitions, and JSON descriptors are cached in RAM on first access.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Connection Pooling&lt;/strong&gt;: SQL and database sinks now reuse shared connection pools across rule actions rather than acquiring new sockets.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Multi-Row SQL Batching&lt;/strong&gt;: For relational sinks (PostgreSQL and SQLite), the sink generates multi-row &lt;code&gt;INSERT INTO ... VALUES (...), (...)&lt;/code&gt; statements in parameterized chunks instead of emitting one query per record.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  Channel Buffer Sizing
&lt;/h3&gt;

&lt;p&gt;We expanded internal actor queue depths from 1,024 to 32,768 records. On a single pinned core, this provides enough buffer headroom to absorb operating system scheduling jitter without propagating backpressure back into the MQTT network loop.&lt;/p&gt;




&lt;h2&gt;
  
  
  3. Finding the Exact Limits: Hierarchical 1k Peak Search
&lt;/h2&gt;

&lt;p&gt;Rather than testing arbitrary rounded rates, we searched for the exact ceiling of each workload using a hierarchical binary ladder:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;10,000 msg/s steps&lt;/strong&gt; to identify the 10k window.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;2,500 msg/s steps&lt;/strong&gt; to narrow the bracket.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;1,000 msg/s steps&lt;/strong&gt; to find the exact tipping point.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  Test Rig (Identical to previous benchmarks)
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Host&lt;/strong&gt;: 12-core x86-64 machine, Docker on WSL2 (cgroup v2).&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Engine Container&lt;/strong&gt;: Pinned to &lt;strong&gt;1 CPU core&lt;/strong&gt;, &lt;strong&gt;1 GiB RAM&lt;/strong&gt;, &lt;code&gt;--memory-swap=1g&lt;/code&gt;, &lt;code&gt;TOKIO_WORKER_THREADS=1&lt;/code&gt;.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Mosquitto Broker&lt;/strong&gt;: Isolated container on separate cores, with a strict &lt;strong&gt;4,096-message / 1 MiB outgoing queue limit&lt;/strong&gt;. If the engine falls behind by even a fraction of a second, the broker drops QoS 0 packets immediately.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Load Generator&lt;/strong&gt;: &lt;code&gt;mqttgen&lt;/code&gt; (our standalone Rust publisher) pushing MQTT 3.1.1 QoS 0 across 8 connections from separate cores.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Verification&lt;/strong&gt;: Exact message-by-message sink validation checking count, unique IDs, per-device aggregates, and zero exceptions.&lt;/li&gt;
&lt;/ul&gt;




&lt;h2&gt;
  
  
  4. The Results: Exact Certified Ceilings
&lt;/h2&gt;

&lt;p&gt;Here are the verified limits for all five workloads in rekuiper v0.500-beta:&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fqnl7lw124yzg0aigipt2.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fqnl7lw124yzg0aigipt2.png" alt=" " width="800" height="442"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Workload&lt;/th&gt;
&lt;th&gt;Scenario&lt;/th&gt;
&lt;th&gt;Certified Ceiling&lt;/th&gt;
&lt;th&gt;First Failure&lt;/th&gt;
&lt;th&gt;Limiting Factor&lt;/th&gt;
&lt;th&gt;Single-Core CPU&lt;/th&gt;
&lt;th&gt;Anon RAM&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W1&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Telemetry Filter (1,000 devices)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;150,000 msg/s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;151,000 msg/s&lt;/td&gt;
&lt;td&gt;CPU saturation (99.4%), broker drops 20.5%&lt;/td&gt;
&lt;td&gt;94.5%&lt;/td&gt;
&lt;td&gt;17.4 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W2&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Device 10s Windows (1,000 devices)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;200,000 msg/s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;210,000 msg/s&lt;/td&gt;
&lt;td&gt;Generator schedule (engine lossless to 240k)&lt;/td&gt;
&lt;td&gt;94.4%&lt;/td&gt;
&lt;td&gt;6.7 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W3&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;ESPHome Topics (10,000 topics, &lt;code&gt;meta(topic)&lt;/code&gt;)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;150,000 msg/s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;151,000 msg/s&lt;/td&gt;
&lt;td&gt;CPU saturation (99.3%), broker drops 5.0%&lt;/td&gt;
&lt;td&gt;97.4%&lt;/td&gt;
&lt;td&gt;16.6 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W4&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Vehicle Windows (10,000 VIN wildcard topics)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;200,000 msg/s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;210,000 msg/s&lt;/td&gt;
&lt;td&gt;Generator schedule (engine lossless to 220k)&lt;/td&gt;
&lt;td&gt;97.5%&lt;/td&gt;
&lt;td&gt;18.1 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W5&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;EV Charger Sessions (&lt;code&gt;SESSIONWINDOW(10, 2)&lt;/code&gt;)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;126,000 msg/s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;127,000 msg/s&lt;/td&gt;
&lt;td&gt;Session drain lag (16s exceeds 5s stability limit)&lt;/td&gt;
&lt;td&gt;86.7%&lt;/td&gt;
&lt;td&gt;6.0 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;




&lt;h2&gt;
  
  
  5. Workload Deep-Dive: What Broke Where
&lt;/h2&gt;

&lt;h3&gt;
  
  
  W1: Telemetry Filter (JSON parsing + condition)
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;SQL&lt;/strong&gt;: &lt;code&gt;SELECT id, device, temp, speed * 3.6 AS speed_kmh FROM telem WHERE temp &amp;gt; 21.0&lt;/code&gt;
&lt;/li&gt;
&lt;li&gt;At &lt;strong&gt;150,000 msg/s&lt;/strong&gt;: 0.00% loss, 94.5% CPU, 17.4 MB RAM, 1.0s drain lag.&lt;/li&gt;
&lt;li&gt;At &lt;strong&gt;151,000 msg/s&lt;/strong&gt;: Single-core CPU hit &lt;strong&gt;99.4%&lt;/strong&gt;. The network thread could not drain the socket fast enough; Mosquitto's 4,096-message queue overflowed and dropped 20.53% of packets.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Verdict&lt;/strong&gt;: 150,000 msg/s is the hard physical CPU ceiling for JSON deserialization, arithmetic projection, and filtering on one core.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  W2: Per-Device 10-second Windows (1,000 devices)
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;SQL&lt;/strong&gt;: &lt;code&gt;SELECT device, count(*) AS n, avg(temp), max(speed) FROM telem GROUP BY device, TUMBLINGWINDOW(ss, 10)&lt;/code&gt;
&lt;/li&gt;
&lt;li&gt;At &lt;strong&gt;200,000 msg/s&lt;/strong&gt;: 0.00% loss, 94.4% CPU, 6.7 MB RAM, 10.0s drain lag.&lt;/li&gt;
&lt;li&gt;At &lt;strong&gt;210,000–240,000 msg/s&lt;/strong&gt;: The engine processed all data with 0.00% loss, but the external generator fell off schedule. At 250,000 msg/s, the pipeline collapsed.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Verdict&lt;/strong&gt;: 200,000 msg/s is our certified on-schedule ceiling. Memory remained at a tiny &lt;strong&gt;6.7 MB&lt;/strong&gt; because aggregations (&lt;code&gt;count&lt;/code&gt;, &lt;code&gt;avg&lt;/code&gt;, &lt;code&gt;max&lt;/code&gt;) update in place without buffering raw rows.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  W3: ESPHome Fleet (10,000 distinct topics, plain text)
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;SQL&lt;/strong&gt;: &lt;code&gt;SELECT meta(topic) AS topic, self AS state FROM telem&lt;/code&gt;
&lt;/li&gt;
&lt;li&gt;At &lt;strong&gt;150,000 msg/s&lt;/strong&gt;: 0.00% loss, 97.4% CPU, 16.6 MB RAM. At the end of sending 4.5 million messages, only 15 messages were in transit.&lt;/li&gt;
&lt;li&gt;At &lt;strong&gt;151,000 msg/s&lt;/strong&gt;: CPU reached &lt;strong&gt;99.3%&lt;/strong&gt;, dropping 5.05% at the broker.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Verdict&lt;/strong&gt;: 150,000 msg/s is the exact ceiling for routing and extracting MQTT metadata across 10,000 dynamic topics.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  W4: Vehicle Wildcard Aggregation (10,000 VIN topics)
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;SQL&lt;/strong&gt;: &lt;code&gt;SELECT device, count(*) AS n, avg(speed), max(temp) FROM telem GROUP BY device, TUMBLINGWINDOW(ss, 10)&lt;/code&gt;
&lt;/li&gt;
&lt;li&gt;Ingests across &lt;code&gt;bench/vehicles/+/telemetry&lt;/code&gt;.&lt;/li&gt;
&lt;li&gt;Sustained &lt;strong&gt;200,000 msg/s&lt;/strong&gt; with 0.00% loss and 18.1 MB RAM.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Verdict&lt;/strong&gt;: Certified at 200,000 msg/s on one core.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  W5: EV Charger Sessions (&lt;code&gt;SESSIONWINDOW&lt;/code&gt;)
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;SQL&lt;/strong&gt;: &lt;code&gt;SELECT device, count(*) AS n, max(speed) FROM telem GROUP BY device, SESSIONWINDOW(ss, 10, 2)&lt;/code&gt;
&lt;/li&gt;
&lt;li&gt;2,000 chargers opening and closing irregular activity sessions.&lt;/li&gt;
&lt;li&gt;Tested in 1,000 msg/s increments:

&lt;ul&gt;
&lt;li&gt;125,000 msg/s: PASS (0.00% loss, 3.0s lag)&lt;/li&gt;
&lt;li&gt;126,000 msg/s: PASS (0.00% loss, 4.0s lag, 86.7% CPU, 6.0 MB RAM)&lt;/li&gt;
&lt;li&gt;127,000 msg/s: FAIL (session close lag spiked to 16.0s, exceeding our 5.0s stability limit).&lt;/li&gt;
&lt;/ul&gt;
&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Verdict&lt;/strong&gt;: Exact certified limit is &lt;strong&gt;126,000 msg/s&lt;/strong&gt;.&lt;/li&gt;
&lt;/ul&gt;




&lt;h2&gt;
  
  
  6. Comparison with Other Engines
&lt;/h2&gt;

&lt;p&gt;Here is the updated head-to-head comparison on the common 5k–100k ladder:&lt;/p&gt;

&lt;h3&gt;
  
  
  Highest Verified Loss-Free Ingest Rate
&lt;/h3&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fgcd21e0y6aq2iuwxhjmw.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fgcd21e0y6aq2iuwxhjmw.png" alt=" " width="800" height="383"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Workload&lt;/th&gt;
&lt;th&gt;rekuiper 0.500&lt;/th&gt;
&lt;th&gt;eKuiper 2.4.1&lt;/th&gt;
&lt;th&gt;Telegraf 1.40.0&lt;/th&gt;
&lt;th&gt;Redpanda Connect 4.109.0&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W1: Telemetry filter&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;150k certified&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;td&gt;50k (backlog)&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W2: Device windows&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;200k certified&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;td&gt;Lost 6–27%&lt;/td&gt;
&lt;td&gt;5k&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W3: ESPHome topics&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;150k certified&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;td&gt;50k (backlog)&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W4: Vehicle windows&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;200k certified&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;td&gt;Inconsistent&lt;/td&gt;
&lt;td&gt;5k&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W5: Charger sessions&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;126k certified&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20k&lt;/td&gt;
&lt;td&gt;Unsupported&lt;/td&gt;
&lt;td&gt;Unsupported&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;h3&gt;
  
  
  Memory at 20,000 msg/s (Engine Anonymous RAM)
&lt;/h3&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Workload&lt;/th&gt;
&lt;th&gt;rekuiper 0.500&lt;/th&gt;
&lt;th&gt;eKuiper 2.4.1&lt;/th&gt;
&lt;th&gt;Telegraf 1.40.0&lt;/th&gt;
&lt;th&gt;Redpanda Connect 4.109.0&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W1: Telemetry filter&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;4.4 MB&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;15 MB&lt;/td&gt;
&lt;td&gt;92 MB&lt;/td&gt;
&lt;td&gt;72 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W2: Device windows&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;6.4 MB&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;536 MB&lt;/td&gt;
&lt;td&gt;52 MB (lossy)&lt;/td&gt;
&lt;td&gt;1,012 MB (crashed)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W3: ESPHome topics&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;4.5 MB&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;43 MB&lt;/td&gt;
&lt;td&gt;85 MB&lt;/td&gt;
&lt;td&gt;68 MB&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W4: Vehicle windows&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;10.2 MB&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;886 MB&lt;/td&gt;
&lt;td&gt;94 MB (lossy)&lt;/td&gt;
&lt;td&gt;993 MB (crashed)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;W5: Charger sessions&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;7.3 MB&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;832 MB&lt;/td&gt;
&lt;td&gt;Unsupported&lt;/td&gt;
&lt;td&gt;Unsupported&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;At 20k msg/s on 10,000 vehicle windows, engines that buffer raw rows in RAM consume &lt;strong&gt;886 MB to 1 GB&lt;/strong&gt;, risking kernel OOM kills on edge hardware. rekuiper maintains &lt;strong&gt;10.2 MB&lt;/strong&gt; by computing aggregates incrementally in fixed-size accumulators.&lt;/p&gt;




&lt;h2&gt;
  
  
  7. Real Sustained Score vs. Fake Buffer Backlog
&lt;/h2&gt;

&lt;p&gt;One crucial lesson from this benchmark audit: &lt;strong&gt;eventual delivery is not throughput.&lt;/strong&gt;&lt;/p&gt;

&lt;p&gt;If an engine accepts 200,000 msg/s for 30 seconds by buffering everything into a huge internal queue, and then spends the next 45 seconds after the publisher stops draining that backlog, that is &lt;strong&gt;not&lt;/strong&gt; a 200k engine. That is an engine surviving on buffer mercy.&lt;/p&gt;

&lt;p&gt;In our benchmark harness:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;The publisher and broker queues are strictly bounded (Mosquitto queue capped at 4,096 messages).&lt;/li&gt;
&lt;li&gt;We record the &lt;strong&gt;upstream source gap&lt;/strong&gt; at the exact millisecond publishing finishes.&lt;/li&gt;
&lt;li&gt;If more than 4,096 messages are backlogged at send-end, or if post-send drain takes longer than 5 seconds, &lt;strong&gt;the trial is marked as a failure&lt;/strong&gt;, even if every single message is eventually written to disk.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;When we claim 150,000 msg/s on W1 and W3, it means the engine processed all 4.5 million messages in real-time with an end-of-send backlog of 15 messages and a 1.0s drain time.&lt;/p&gt;




&lt;h2&gt;
  
  
  Trying It Out
&lt;/h2&gt;

&lt;p&gt;rekuiper is free and open source under &lt;strong&gt;MIT / Apache-2.0&lt;/strong&gt;:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;GitHub Repository&lt;/strong&gt;: &lt;a href="https://github.com/ankur-paan/rekuiper" rel="noopener noreferrer"&gt;github.com/ankur-paan/rekuiper&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Prebuilt Binaries&lt;/strong&gt;: Available for Linux (x86_64), macOS (Intel &amp;amp; Apple Silicon), and Windows under &lt;a href="https://github.com/ankur-paan/rekuiper/releases" rel="noopener noreferrer"&gt;Releases v0.500-beta&lt;/a&gt;.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Docker&lt;/strong&gt;:
&lt;/li&gt;
&lt;/ul&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;  docker run &lt;span class="nt"&gt;-d&lt;/span&gt; &lt;span class="nt"&gt;--name&lt;/span&gt; rekuiper &lt;span class="se"&gt;\&lt;/span&gt;
    &lt;span class="nt"&gt;-p&lt;/span&gt; 9081:9081 &lt;span class="nt"&gt;-p&lt;/span&gt; 20499:20499 &lt;span class="se"&gt;\&lt;/span&gt;
    &lt;span class="nt"&gt;-e&lt;/span&gt; &lt;span class="nv"&gt;KUIPER__BASIC__CONSOLELOG&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="nb"&gt;true&lt;/span&gt; &lt;span class="se"&gt;\&lt;/span&gt;
    ankurkrp/rekuiper:0.500-beta
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;All raw evidence files, harness scripts (&lt;code&gt;mqttgen&lt;/code&gt;, &lt;code&gt;iotrunner&lt;/code&gt;), and reproduction instructions are published in &lt;a href="https://github.com/ankur-paan/rekuiper/tree/main/test/benchmark/iiot-mqtt" rel="noopener noreferrer"&gt;&lt;code&gt;test/benchmark/iiot-mqtt/&lt;/code&gt;&lt;/a&gt;.&lt;/p&gt;

</description>
      <category>rust</category>
      <category>iot</category>
      <category>performance</category>
      <category>opensource</category>
    </item>
    <item>
      <title>At the edge, the number that matters is memory - not throughput (specially in Ramageddon)</title>
      <dc:creator>Ankur Kumar Pandey</dc:creator>
      <pubDate>Sun, 13 Sep 2026 22:43:55 +0000</pubDate>
      <link>https://dev.to/ankurpaan/at-the-edge-the-number-that-matters-is-memory-not-throughput-specially-in-ramageddon-5hc7</link>
      <guid>https://dev.to/ankurpaan/at-the-edge-the-number-that-matters-is-memory-not-throughput-specially-in-ramageddon-5hc7</guid>
      <description>&lt;h3&gt;
  
  
  We rebuilt LF Edge eKuiper in Rust and ran it against eKuiper, Telegraf and Redpanda Connect on five real MQTT workloads — one core, 1 GB of memory, and output checked message-by-message. The most important result wasn't speed.
&lt;/h3&gt;

&lt;p&gt;&lt;em&gt;I-Dacs Labs Engineering · ~16 min read&lt;/em&gt;&lt;/p&gt;




&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fdzcvqqtcg581bdyebs2h.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fdzcvqqtcg581bdyebs2h.png" alt=" " width="799" height="420"&gt;&lt;/a&gt;&lt;/p&gt;




&lt;p&gt;Most stream-processing benchmarks you'll read optimize for one number: peak throughput on a big server. That number is close to useless for the place these engines actually run — an industrial gateway, an ESPHome hub, a vehicle head-unit, an EV charger. There, you get one or two CPU cores and a few hundred megabytes of free memory, your input arrives over MQTT, and your traffic is bursty in the worst way: fleets reconnect together, chargers start sessions together, devices flush buffered readings all at once after an outage.&lt;/p&gt;

&lt;p&gt;In that world two questions decide whether your pipeline survives, and neither is peak throughput:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;Does the engine keep up &lt;strong&gt;on a single core&lt;/strong&gt;?&lt;/li&gt;
&lt;li&gt;Does its memory &lt;strong&gt;stay bounded&lt;/strong&gt; when traffic grows?&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;We built a stream engine called &lt;strong&gt;rekuiper&lt;/strong&gt; to answer "yes" to both, and then we built a benchmark honest enough to tell us if we'd actually managed it. This post is about the benchmark as much as the engine, because the benchmark taught us more than we expected — including a correctness bug in our own code that a throughput-only test would have rewarded as "fast."&lt;/p&gt;

&lt;p&gt;The headline: across five MQTT workloads shaped like real deployments, rekuiper produced complete, correct output at &lt;strong&gt;100,000 messages per second on one core&lt;/strong&gt; in every workload — the top of our tested range, so we never found its ceiling. But the result we care about most isn't that. It's that on the windowed workloads, rekuiper's memory stayed between &lt;strong&gt;5 and 10 MB&lt;/strong&gt; while the Go-based engines climbed to &lt;strong&gt;half a gigabyte to a full gigabyte&lt;/strong&gt;, or failed. That gap is the whole point, and it comes from design, not from the language.&lt;/p&gt;




&lt;h2&gt;
  
  
  What rekuiper is
&lt;/h2&gt;

&lt;p&gt;rekuiper is a stream-processing engine written in Rust that reimplements the surface of &lt;strong&gt;LF Edge eKuiper&lt;/strong&gt;: its REST API, its SQL dialect, its stream and rule definitions, and its &lt;code&gt;kuiper&lt;/code&gt; command-line interface. The goal was boring on purpose — existing eKuiper rules, the eKuiper Manager web UI, and deployment tooling should keep working — so that "switch the engine" isn't also "rewrite everything."&lt;/p&gt;

&lt;p&gt;Concretely, the compatibility surface covers eKuiper's REST API (98 paths and 140 operations, checked black-box against eKuiper's own OpenAPI description), the SQL dialect including JSON paths, &lt;code&gt;CASE&lt;/code&gt;, array indexing and &lt;code&gt;unnest&lt;/code&gt;, and eKuiper's stream option names (&lt;code&gt;DATASOURCE&lt;/code&gt;, &lt;code&gt;FORMAT&lt;/code&gt;, &lt;code&gt;CONF_KEY&lt;/code&gt;, &lt;code&gt;SCHEMAID&lt;/code&gt;, &lt;code&gt;TIMESTAMP&lt;/code&gt;, and so on). If you know eKuiper, you already know rekuiper.&lt;/p&gt;

&lt;p&gt;What's different is underneath, and it's built around one principle: &lt;strong&gt;memory stays bounded under load.&lt;/strong&gt; Three design choices carry that, and each one shows up later in the numbers.&lt;/p&gt;




&lt;h2&gt;
  
  
  Three design choices that keep memory flat
&lt;/h2&gt;

&lt;h3&gt;
  
  
  Bounded queues with real backpressure
&lt;/h3&gt;

&lt;p&gt;Sources publish records into an in-process &lt;strong&gt;stream bus&lt;/strong&gt; with bounded per-subscriber queues — 4,096 records each. Admission is &lt;strong&gt;reserve-then-commit&lt;/strong&gt;: a batch first reserves capacity in &lt;em&gt;every&lt;/em&gt; subscriber's queue, and only then commits. So a batch is either delivered to all subscribers or rejected outright, never half-delivered, and a slow rule pushes back on its source instead of quietly dropping data. Each rule runs as its own task, and its output drains through a bounded sink queue (default 10,000) served by a dedicated sink worker.&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2F4th6cl5jpfx6u2jksoii.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2F4th6cl5jpfx6u2jksoii.png" alt=" " width="799" height="521"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;p&gt;The MQTT source uses the &lt;code&gt;rumqttc&lt;/code&gt; client, and when a single network read surfaces several publishes, the source admits everything already buffered as &lt;strong&gt;one batch of up to 1,024 records&lt;/strong&gt;. That avoids a per-message wakeup without ever waiting around for more data — you pay one scheduling cost for a burst instead of one per message. This batch-admission trick is a big part of why rekuiper uses roughly half the CPU per message of the Go engines on the simple workloads.&lt;/p&gt;

&lt;h3&gt;
  
  
  Incremental window aggregation: O(groups), not O(messages)
&lt;/h3&gt;

&lt;p&gt;This is the important one. When you compute &lt;code&gt;GROUP BY device, TUMBLINGWINDOW(ss, 10)&lt;/code&gt; with &lt;code&gt;count&lt;/code&gt;, &lt;code&gt;avg&lt;/code&gt;, &lt;code&gt;max&lt;/code&gt; and friends, the naive way is to buffer every row that falls in the window and aggregate at the trigger. Memory then grows with &lt;strong&gt;traffic&lt;/strong&gt; — messages per window — which is exactly the thing that explodes when a fleet reconnects.&lt;/p&gt;

&lt;p&gt;rekuiper instead keeps &lt;strong&gt;one accumulator per group per aggregate&lt;/strong&gt; and never stores the rows. Window memory becomes a function of the &lt;strong&gt;number of devices&lt;/strong&gt;, not the number of messages. For the common edge shape — group columns, plain columns, and &lt;code&gt;count&lt;/code&gt;/&lt;code&gt;sum&lt;/code&gt;/&lt;code&gt;avg&lt;/code&gt;/&lt;code&gt;min&lt;/code&gt;/&lt;code&gt;max&lt;/code&gt; over simple expressions — this incremental evaluator does the whole job. Statements that genuinely need the rows (&lt;code&gt;collect()&lt;/code&gt;, joins, some &lt;code&gt;HAVING&lt;/code&gt;) fall back to a buffered evaluator, and a unit test checks the two produce identical output on mixed data.&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fhmgosicln856s905alym.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fhmgosicln856s905alym.png" alt=" " width="800" height="457"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;p&gt;For a fleet, this is the difference between memory that scales with how many vehicles you have and memory that scales with how fast they're all talking at once. Only one of those is safe on a 1 GB box.&lt;/p&gt;

&lt;h3&gt;
  
  
  An offline sink cache that spills instead of dropping
&lt;/h3&gt;

&lt;p&gt;For intermittent uplinks — a vehicle in a tunnel, a remote site on flaky cellular — a sink can enable a cache using eKuiper's own options (&lt;code&gt;enableCache&lt;/code&gt;, &lt;code&gt;memoryCacheThreshold&lt;/code&gt;, &lt;code&gt;maxDiskCache&lt;/code&gt;, and the rest). Records whose send fails recoverably are queued FIFO: in memory up to a threshold, then in disk pages, and only when the disk budget is exhausted are the oldest records dropped — and counted, not silently lost. The MQTT sink holds one persistent connection per action and reports disconnection, so an outage is detected and cached rather than quietly discarded. (The cache is covered by an integration test but isn't part of the performance numbers here.)&lt;/p&gt;




&lt;h2&gt;
  
  
  The benchmark that doesn't lie to you
&lt;/h2&gt;

&lt;p&gt;Here's the uncomfortable truth about a lot of edge stream-processing comparisons: they measure throughput at the point the engine &lt;em&gt;acknowledges&lt;/em&gt; ingest, or they count output records without checking that the records are &lt;em&gt;correct&lt;/em&gt;. Both can hide loss and duplication completely. An engine that drops 15% of your data can look fast if you never verify what came out the other end.&lt;/p&gt;

&lt;p&gt;So we built the benchmark around &lt;strong&gt;exact output verification&lt;/strong&gt;, and gave every engine the same cramped room to work in.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Equal, realistic limits.&lt;/strong&gt; Every engine runs in a container pinned to &lt;strong&gt;one CPU core with 1 GB of memory and no swap&lt;/strong&gt; (&lt;code&gt;--cpuset-cpus=2 --cpus=1 --memory=1g --memory-swap=1g&lt;/code&gt;). A separate Mosquitto broker gets its own cores and generous queue limits, so the broker is never the bottleneck. An open-loop Rust load generator (&lt;code&gt;mqttgen&lt;/code&gt;, standard library only, MQTT 3.1.1, QoS 0) feeds every engine from the same schedule, and a step only counts if the generator actually stayed on schedule.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Four engines.&lt;/strong&gt; rekuiper v0.425-beta, &lt;strong&gt;eKuiper 2.4.1&lt;/strong&gt;, &lt;strong&gt;Telegraf 1.40.0&lt;/strong&gt;, and &lt;strong&gt;Redpanda Connect 4.109.0&lt;/strong&gt; (formerly Benthos). We deliberately &lt;em&gt;excluded&lt;/em&gt; Apache Flink: neither Flink 2.x nor Apache Bahir ships an MQTT connector, so testing Flink would have meant a custom source or a Kafka bridge — changing the very ingest path under test. Rather than benchmark a different pipeline and call it Flink, we left it out and said so.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Five workloads shaped like real deployments:&lt;/strong&gt;&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;W1 — telemetry filter.&lt;/strong&gt; 1,000 devices, one topic, a simple &lt;code&gt;WHERE temp &amp;gt; 21.0&lt;/code&gt; with a unit conversion. Stateless.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;W2 — per-device windows.&lt;/strong&gt; 1,000 devices, 10-second tumbling windows with &lt;code&gt;count&lt;/code&gt;/&lt;code&gt;avg&lt;/code&gt;/&lt;code&gt;max&lt;/code&gt;. Stateful.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;W3 — ESPHome states.&lt;/strong&gt; 10,000 plain-text topics via wildcard, using &lt;code&gt;FORMAT="binary"&lt;/code&gt; and &lt;code&gt;meta(topic)&lt;/code&gt; to carry the topic through. Stateless but wide.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;W4 — vehicle windows.&lt;/strong&gt; 10,000 topics (one per VIN), 10-second tumbling windows. Stateful and wide — the hardest memory test.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;W5 — EV charger sessions.&lt;/strong&gt; 2,000 topics, &lt;code&gt;SESSIONWINDOW(ss, 10, 2)&lt;/code&gt;. Neither Telegraf nor Redpanda Connect has a session window, so they can't express it at all.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;&lt;strong&gt;Proofs, not vibes.&lt;/strong&gt; For each engine, workload and rate (5k, 20k, 50k, 100k msg/s), we warm up until the subscription is provably live, send for 30 seconds on a fixed schedule, drain until the sink file stops growing, then verify the output &lt;em&gt;exactly&lt;/em&gt;:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;W1:&lt;/strong&gt; the count of unique message IDs carrying the run tag must equal the closed-form expected filtered count, with &lt;strong&gt;no duplicates&lt;/strong&gt;.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;W2, W4, W5:&lt;/strong&gt; the sum of per-device counts across all output windows must equal the messages sent, and &lt;strong&gt;every device must appear&lt;/strong&gt;.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;W3:&lt;/strong&gt; output rows must equal messages sent, and &lt;strong&gt;all 10,000 topics&lt;/strong&gt; must appear.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;A step is &lt;strong&gt;complete&lt;/strong&gt; only when its proof holds. "Loss" is the relative shortfall against the proof. This is the part that makes the numbers trustworthy — and, as you'll see, it's the part that caught our own bug.&lt;/p&gt;




&lt;h2&gt;
  
  
  Results
&lt;/h2&gt;

&lt;h3&gt;
  
  
  Nobody else finished the range
&lt;/h3&gt;

&lt;p&gt;&lt;strong&gt;Highest tested rate with complete, correct output:&lt;/strong&gt;&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Foaz58bau44twpku7wrb6.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Foaz58bau44twpku7wrb6.png" alt=" " width="800" height="412"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Workload&lt;/th&gt;
&lt;th&gt;rekuiper&lt;/th&gt;
&lt;th&gt;eKuiper 2.4.1&lt;/th&gt;
&lt;th&gt;Telegraf 1.40.0&lt;/th&gt;
&lt;th&gt;Redpanda Connect 4.109.0&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;W1 telemetry filter&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;≥ 100,000&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;td&gt;50,000 (lag 10 s)&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W2 per-device windows&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;≥ 100,000&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;td&gt;none&lt;/td&gt;
&lt;td&gt;5,000&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W3 ESPHome states&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;≥ 100,000&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;td&gt;50,000 (lag 5 s)&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W4 vehicle windows&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;≥ 100,000&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;td&gt;50,000 only&lt;/td&gt;
&lt;td&gt;5,000&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W5 charger sessions&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;≥ 100,000&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;20,000&lt;/td&gt;
&lt;td&gt;not supported&lt;/td&gt;
&lt;td&gt;not supported&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;rekuiper completed all 20 steps. Because 100,000 msg/s was the top of the range, we never reached its limit — at 100k on the wide ESPHome workload it used 96.6% of the core, and the windowed workloads used 78–85%, so there's headroom left. eKuiper was solid and complete up to 20,000 msg/s across the board. Telegraf managed 50,000 on two stateless workloads but never produced complete per-device windows at any rate. Redpanda Connect reached 20,000 on stateless workloads and 5,000 on windows.&lt;/p&gt;

&lt;h3&gt;
  
  
  The memory gap
&lt;/h3&gt;

&lt;p&gt;This is the result we'd frame and put on the wall. Peak engine heap (cgroup anonymous memory) at &lt;strong&gt;20,000 msg/s&lt;/strong&gt;, the highest rate every engine could still be compared at:&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2F49jezc3jw5gm8oizo26a.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2F49jezc3jw5gm8oizo26a.png" alt=" " width="800" height="469"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Workload&lt;/th&gt;
&lt;th&gt;rekuiper&lt;/th&gt;
&lt;th&gt;eKuiper&lt;/th&gt;
&lt;th&gt;Telegraf&lt;/th&gt;
&lt;th&gt;Redpanda Connect&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;W1&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;4.7&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;15.3&lt;/td&gt;
&lt;td&gt;91.6&lt;/td&gt;
&lt;td&gt;71.5&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W2&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;5.4&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;535.9&lt;/td&gt;
&lt;td&gt;51.6 †&lt;/td&gt;
&lt;td&gt;1,012.3 †&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W3&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;4.8&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;43.2&lt;/td&gt;
&lt;td&gt;84.9&lt;/td&gt;
&lt;td&gt;67.7&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W4&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;10.0&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;886.2&lt;/td&gt;
&lt;td&gt;94.0 †&lt;/td&gt;
&lt;td&gt;992.9 †&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;W5&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;5.9&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;831.7&lt;/td&gt;
&lt;td&gt;n/a&lt;/td&gt;
&lt;td&gt;n/a&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;&lt;em&gt;(† marks a step whose correctness proof failed — the memory figure is real, but the engine wasn't producing complete output.)&lt;/em&gt;&lt;/p&gt;

&lt;p&gt;On the windowed workloads (W2, W4, W5), eKuiper's heap ran to &lt;strong&gt;536–886 MB&lt;/strong&gt; and Redpanda Connect's &lt;code&gt;system_window&lt;/code&gt; pattern — which holds every message of a window before aggregating — hit the &lt;strong&gt;1 GB&lt;/strong&gt; ceiling. rekuiper stayed at &lt;strong&gt;5–10 MB&lt;/strong&gt; at every rate on every workload. Two orders of magnitude, on the exact workload edge fleets generate.&lt;/p&gt;

&lt;h3&gt;
  
  
  CPU: roughly half
&lt;/h3&gt;

&lt;p&gt;At 20,000 msg/s, rekuiper used &lt;strong&gt;44–49% of one core&lt;/strong&gt;. eKuiper used &lt;strong&gt;86–99%&lt;/strong&gt;, and the two Go pipeline tools were similar or worse (on the steps where they were even producing correct output). About half the CPU per message, which on a shared single-core box is the difference between comfortable headroom and being one traffic spike away from falling behind.&lt;/p&gt;




&lt;h2&gt;
  
  
  Where the difference actually comes from
&lt;/h2&gt;

&lt;p&gt;It would be easy, and wrong, to write this up as "Rust beats Go." The language helps, but the CPU difference on the stateless workloads is roughly &lt;strong&gt;2×, not 10×&lt;/strong&gt;, because MQTT receive, JSON decode and file writing dominate for everyone. The interesting gaps are structural.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Window memory is a design choice, not a language feature.&lt;/strong&gt; eKuiper's windowed memory grows with message rate — about 125 MB at 5,000 msg/s, 536–886 MB at 20,000 — and hits the 1 GB limit at 50,000, where output loss immediately follows. Redpanda Connect's documented windowing buffers the whole window and reaches the limit from 20,000 msg/s. rekuiper's incremental evaluator keeps one accumulator per device, so heap is a function of &lt;strong&gt;fleet size&lt;/strong&gt;, not traffic. Any of these engines could adopt the same approach; the point is that it's the approach, not the runtime, that matters here.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;The loss mechanism is the broker, honestly reported.&lt;/strong&gt; When an engine falls behind, its MQTT subscription backs up and the broker drops QoS 0 messages for that slow subscriber. This is visible directly in the input counters — for instance, eKuiper received only 1.23 of 3.0 million messages on W1 at 100,000 msg/s. Under QoS 1 the same overload would surface as backpressure on the publishers instead of loss. We report QoS 0 because it's the common, cheap edge default, and because it makes overload measurable rather than hidden.&lt;/p&gt;




&lt;h2&gt;
  
  
  The bug our own benchmark caught
&lt;/h2&gt;

&lt;p&gt;Here's the part we could have quietly left out, and won't.&lt;/p&gt;

&lt;p&gt;An earlier run of this exact harness, on a previous rekuiper build, showed about &lt;strong&gt;15% "loss" on W2 at every rate&lt;/strong&gt;, and 100% loss with 820 MB of memory at 100,000 msg/s. It looked like overload. It wasn't — it was a correctness defect in our window evaluation.&lt;/p&gt;

&lt;p&gt;Our time-window trigger was collapsing the whole window into a &lt;strong&gt;single aggregate&lt;/strong&gt;: it ignored &lt;code&gt;GROUP BY&lt;/code&gt; partitioning (emitting one row per window, with group values taken from the first record), ignored &lt;code&gt;WHERE&lt;/code&gt;, and buffered and cloned every row on the way. The reason it produced a suspiciously &lt;em&gt;constant&lt;/em&gt; shortfall was subtle: the warm-up device's row was absorbing the first window of measured data every time.&lt;/p&gt;

&lt;p&gt;Our unit tests didn't catch it, because they aggregated a single group — exactly the case the bug handled correctly. Only the &lt;strong&gt;exact per-device proof&lt;/strong&gt; in the benchmark exposed it. We fixed it (that fix is the incremental evaluator described above) before the final measurements.&lt;/p&gt;

&lt;p&gt;We're telling you this because it's the strongest argument in the whole paper for verifying output content: a throughput-only benchmark, or one that counts output rows without checking their identity, would have looked at that defective build and reported it as fast and lean. The bug &lt;em&gt;reduced&lt;/em&gt; work by skipping grouping and filtering. Speed without a correctness proof is not a measurement; it's a guess with a stopwatch.&lt;/p&gt;




&lt;h2&gt;
  
  
  A "neutral" harness detail that reordered the results
&lt;/h2&gt;

&lt;p&gt;One more methodology lesson, because it surprised us.&lt;/p&gt;

&lt;p&gt;In an earlier comparison, the output sink file lived on a Windows-drive bind mount, where every write system call is expensive. Telegraf's file output, by default, issues one unbuffered write per metric; Redpanda Connect writes each message individually. Both were &lt;strong&gt;bound by the sink, not their own logic&lt;/strong&gt; — flat CPU around 40–47% while losing up to 89% of messages — purely because of where the file lived. eKuiper was affected too, though less. rekuiper batches its file writes, so it barely noticed.&lt;/p&gt;

&lt;p&gt;Moving the sink to the VM's local ext4 filesystem changed the standings substantially. eKuiper W1 at 20,000 msg/s went from 1.7% loss to complete; Telegraf W1 at 20,000 went from 37% loss to complete. Same engines, same rates, same rules — different disk. We kept the superseded runs in the artifact and marked them as such, and the lesson is now a rule we'd apply to any stream-processor benchmark: &lt;strong&gt;state where your sink writes and how often it issues system calls&lt;/strong&gt;, because that detail can quietly decide your rankings.&lt;/p&gt;

&lt;p&gt;(There's an honest loose end here too: Telegraf lost a near-constant 9.4% on the windowed workloads at low rates even with a grace period, yet was complete at 50,000 on W4. We didn't find the cause. The configuration is published so someone else can.)&lt;/p&gt;




&lt;h2&gt;
  
  
  What this doesn't prove
&lt;/h2&gt;

&lt;p&gt;We build rekuiper, so treat the framing with the skepticism it deserves — and here's what to hold against it:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;One host, one repetition.&lt;/strong&gt; The whole run was a single laptop under Windows 11 and WSL2, one repetition per cell. Host-health snapshots flag several runs as noisy. The gaps are large relative to that noise, and an earlier run showed the same qualitative pattern, but repeated runs on a dedicated Linux host would be stronger.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;A coarse rate ladder.&lt;/strong&gt; "20,000" means an engine passed 20,000 and failed at 50,000; the true limit is somewhere in between. rekuiper's ceiling wasn't measured at all.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Defaults, mostly.&lt;/strong&gt; eKuiper ran with default rule options; a partial run with larger buffers showed a similar pattern but wasn't repeated. Redpanda Connect used its documented windowing pattern — other designs might do better.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;QoS 0 only.&lt;/strong&gt; No QoS 1/2, TLS, broker reconnect storms, event-time or out-of-order windows, joins, or non-MQTT connectors. rekuiper's MQTT path is its stable, benchmarked path; its other connectors are not yet considered stable.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;None of that changes the central, two-order-of-magnitude memory result, but you should know exactly where the edges of the claim are.&lt;/p&gt;




&lt;h2&gt;
  
  
  Try it, and break it
&lt;/h2&gt;

&lt;p&gt;Everything here is reproducible. The engine, the orchestrator, the load generator, every configuration, and the raw per-step evidence (per-second CPU and memory, generator reports, input counters, host health, image IDs) are published at tag &lt;code&gt;v0.425-beta&lt;/code&gt;:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;git clone https://github.com/ankur-paan/rekuiper.git
&lt;span class="nb"&gt;cd &lt;/span&gt;rekuiper
&lt;span class="c"&gt;# Method, configs and commands:&lt;/span&gt;
&lt;span class="c"&gt;#   test/benchmark/iiot-mqtt/README.md&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;A full run — four engines, five workloads, four rates — takes about 1 hour 45 minutes and needs Linux or WSL2 with Docker (cgroup v2), a Rust toolchain, and at least 12 logical CPUs.&lt;/p&gt;

&lt;p&gt;If you run IIoT gateways, ESPHome fleets, vehicle telemetry, or EV chargers, the most useful thing you can do with this is try to break it on your own hardware and your own rules — especially ARM, and especially with buffer settings tuned for your traffic. We'd genuinely rather hear where it falls over than where it wins.&lt;/p&gt;

&lt;p&gt;Because at the edge, the engine that survives a Monday-morning reconnect storm isn't the one with the biggest throughput number. It's the one whose memory you can still predict when ten thousand devices all start talking at once.&lt;/p&gt;

&lt;p&gt;👉 &lt;strong&gt;GitHub:&lt;/strong&gt; &lt;a href="https://github.com/ankur-paan/rekuiper" rel="noopener noreferrer"&gt;https://github.com/ankur-paan/rekuiper&lt;/a&gt;&lt;/p&gt;




&lt;p&gt;&lt;em&gt;rekuiper v0.425-beta is dual-licensed MIT / Apache-2.0.&lt;/em&gt;&lt;/p&gt;

</description>
      <category>rust</category>
      <category>iot</category>
      <category>performance</category>
      <category>opensource</category>
    </item>
    <item>
      <title>We rewrote LF Edge eKuiper in Rust: 425k eps, 8MB RAM, 13ms boot (vs Flink, Go eKuiper &amp; Benthos)</title>
      <dc:creator>Ankur Kumar Pandey</dc:creator>
      <pubDate>Thu, 10 Sep 2026 03:30:28 +0000</pubDate>
      <link>https://dev.to/ankurpaan/we-rewrote-lf-edge-ekuiper-in-rust-425k-eps-8mb-ram-13ms-boot-vs-flink-go-ekuiper-benthos-1o3p</link>
      <guid>https://dev.to/ankurpaan/we-rewrote-lf-edge-ekuiper-in-rust-425k-eps-8mb-ram-13ms-boot-vs-flink-go-ekuiper-benthos-1o3p</guid>
      <description>&lt;p&gt;Over at &lt;strong&gt;&lt;a href="https://i-dacs.com" rel="noopener noreferrer"&gt;I-Dacs Labs&lt;/a&gt;&lt;/strong&gt;, we run high-throughput telemetry pipelines on edge devices (Raspberry Pis, Advantech gateways, and embedded x86/ARM boxes). We've been using LF Edge eKuiper for local stream processing (SQL filtering, sliding windows, and MQTT/Kafka sinks), but kept hitting the classic edge computing wall:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;JVM engines (Apache Flink)&lt;/strong&gt;: Incredible throughput, but require &amp;gt;1 GB RAM and take 20+ seconds to boot. Unusable on small industrial hardware.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Go engines (Upstream eKuiper, Benthos, Telegraf)&lt;/strong&gt;: Much lighter, but continuous Stop-The-World GC sweeps introduced tail latency jitter. Worse, under burst sensor loads (10k–100k events/sec), Go channel buffer saturation led to silent packet loss.
We decided to rewrite the entire engine in pure Rust: &lt;strong&gt;&lt;a href="https://github.com/ankur-paan/rekuiper" rel="noopener noreferrer"&gt;rekuiper&lt;/a&gt;&lt;/strong&gt; (v0.421-beta, dual licensed under MIT / Apache-2.0).
---
### What We Built&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Core Architecture&lt;/strong&gt;: Lock-free stream bus (&lt;code&gt;StreamBus&lt;/code&gt;), Tokio async actors for rule execution, and bounded actor queues for sinks with zero runtime GC pauses.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;100% Drop-In Parity&lt;/strong&gt;: Fully compatible with the existing eKuiper Manager Web UI, OpenAPI 3.0 schemas, and standard streaming SQL. Zero scaffolded stubs across all 98 REST endpoints.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Footprint&lt;/strong&gt;: 9.60 MB stripped static binary, ~6 – 8.2 MB idle RAM consumption.&lt;/li&gt;
&lt;/ol&gt;

&lt;h2&gt;
  
  
  - &lt;strong&gt;Sub-15ms Cold Boot&lt;/strong&gt;: 12.5 – 14.5 ms internal daemon bootstrap; 123 ms end-to-end process-spawn-to-ready.
&lt;/h2&gt;

&lt;h3&gt;
  
  
  Empirical Head-to-Head Benchmarks (500,000 Records)
&lt;/h3&gt;

&lt;p&gt;Rather than hand-waving estimates, we ran all engines head-to-head on the exact same Linux machine (WSL2 / Ubuntu x86_64) using an identical 500,000-record telemetry workload:&lt;br&gt;
&lt;strong&gt;Pipeline:&lt;/strong&gt; Parse 500k JSON events → Compute formula (&lt;code&gt;temp * 1.8 + 32&lt;/code&gt;) → Filter (&lt;code&gt;temp &amp;gt; 20.0&lt;/code&gt;) → Project (&lt;code&gt;id, temp_f&lt;/code&gt;) → Sink&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
  &lt;thead&gt;
    &lt;tr&gt;
      &lt;th&gt;Engine&lt;/th&gt;
      &lt;th&gt;Runtime&lt;/th&gt;
      &lt;th&gt;500k Elapsed&lt;/th&gt;
      &lt;th&gt;Throughput&lt;/th&gt;
      &lt;th&gt;Data Drops&lt;/th&gt;
      &lt;th&gt;Memory&lt;/th&gt;
    &lt;/tr&gt;
  &lt;/thead&gt;
  &lt;tbody&gt;
    &lt;tr&gt;
      &lt;td&gt;&lt;strong&gt;rekuiper (0.421)&lt;/strong&gt;&lt;/td&gt;
      &lt;td&gt;&lt;strong&gt;Pure Rust&lt;/strong&gt;&lt;/td&gt;
      &lt;td&gt;&lt;strong&gt;1.176 s&lt;/strong&gt;&lt;/td&gt;
      &lt;td&gt;&lt;strong&gt;425,308 eps&lt;/strong&gt;&lt;/td&gt;
      &lt;td&gt;&lt;strong&gt;0 (0.0%)&lt;/strong&gt;&lt;/td&gt;
      &lt;td&gt;&lt;strong&gt;~8 MB&lt;/strong&gt;&lt;/td&gt;
    &lt;/tr&gt;
    &lt;tr&gt;
      &lt;td&gt;Apache Flink&lt;/td&gt;
      &lt;td&gt;Java / JVM&lt;/td&gt;
      &lt;td&gt;2.144 s&lt;/td&gt;
      &lt;td&gt;233,209 eps&lt;/td&gt;
      &lt;td&gt;0 (0.0%)&lt;/td&gt;
      &lt;td&gt;~1,022 MB&lt;/td&gt;
    &lt;/tr&gt;
    &lt;tr&gt;
      &lt;td&gt;Telegraf&lt;/td&gt;
      &lt;td&gt;Go&lt;/td&gt;
      &lt;td&gt;8.194 s&lt;/td&gt;
      &lt;td&gt;61,019 eps&lt;/td&gt;
      &lt;td&gt;0 (0.0%)&lt;/td&gt;
      &lt;td&gt;~50 MB&lt;/td&gt;
    &lt;/tr&gt;
    &lt;tr&gt;
      &lt;td&gt;Upstream Go eKuiper&lt;/td&gt;
      &lt;td&gt;Go&lt;/td&gt;
      &lt;td&gt;11.290 s&lt;/td&gt;
      &lt;td&gt;44,287 eps&lt;/td&gt;
      &lt;td&gt;72,921 (14.6%)&lt;/td&gt;
      &lt;td&gt;~45 MB&lt;/td&gt;
    &lt;/tr&gt;
    &lt;tr&gt;
      &lt;td&gt;Redpanda Connect&lt;/td&gt;
      &lt;td&gt;Go&lt;/td&gt;
      &lt;td&gt;19.236 s&lt;/td&gt;
      &lt;td&gt;25,993 eps&lt;/td&gt;
      &lt;td&gt;0 (0.0%)&lt;/td&gt;
      &lt;td&gt;~38 MB&lt;/td&gt;
    &lt;/tr&gt;
  &lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;


&lt;h3&gt;
  
  
  Key Observations:
&lt;/h3&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;Channel Saturation in Go&lt;/strong&gt;: Under sustained 500k burst ingestion, upstream Go eKuiper dropped &lt;strong&gt;72,921 records&lt;/strong&gt; (14.6% data loss) due to channel saturation (&lt;code&gt;buffer full, drop message&lt;/code&gt;). &lt;code&gt;rekuiper&lt;/code&gt; processed all 500,000 events with 0 drops in 1.176s (&lt;strong&gt;9.6x faster&lt;/strong&gt;).&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;vs Apache Flink&lt;/strong&gt;: Flink’s execution graph is fast (233k eps), but the JobManager + TaskManager JVM consumed &lt;strong&gt;over 1 GB of RAM&lt;/strong&gt;. &lt;code&gt;rekuiper&lt;/code&gt; beats it in single-core throughput while consuming &lt;strong&gt;125x less RAM&lt;/strong&gt; (&amp;lt; 8.2 MB).&lt;/li&gt;
&lt;/ol&gt;
&lt;h2&gt;
  
  
  3. &lt;strong&gt;Cold Boot Time&lt;/strong&gt;: &lt;code&gt;rekuiper&lt;/code&gt; boots internally in &lt;strong&gt;~13 ms&lt;/strong&gt; (123 ms OS spawn to socket ready), compared to 1.2s for Go eKuiper and 20s for Apache Flink.
&lt;/h2&gt;
&lt;h3&gt;
  
  
  Reproduce It in 2 Minutes
&lt;/h3&gt;

&lt;p&gt;All reproduction scripts and Docker configs are in the repository. Anyone can run the whole suite:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;git clone https://github.com/ankur-paan/rekuiper.git
&lt;span class="nb"&gt;cd &lt;/span&gt;rekuiper
./test/benchmark/run_all.sh
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



</description>
      <category>rust</category>
      <category>iot</category>
      <category>performance</category>
      <category>opensource</category>
    </item>
  </channel>
</rss>
