1. Introduction
The data pipelines come in different shapes such as streaming telemetry data, machine learning, messaging, or other real-time data feed. Common solutions to these problems include messaging middlewares like Kafka, Apache Flink/Spark, Airflow DAGs, Temporal, CDC via Debezium, etc. These tools may differ, but the design underneath is the same: ingest data, pipe it through a chain of stages that may hold state, and then send it to the destination after some filtering/decoration/massaging of the data. This post argues that this kind of stateful pipe and filter based architecture can be build with a simple primitive: durable actor connected by channels, which is what I built PlexSpaces around. This architecture is also similar to the SEDA (Staged Event-Driven Architecture) that builds applications as a network of stages joined by event queues. Each stage has its own workers with a queue that give you backpressure, isolation, and a place to observe load, e.g, :

An actor is a SEDA stage in miniature:

By default actor mailbox is in-memory state and it doesn’t connect a stage to Kafka, SQS, etc. PlexSpaces addresses these gaps with:
- Durable actors. The durability facet journals state changes and checkpoints.
- Channels. Queue-based coupling between stages.
In addition, PlexSpaces supports APIs for TupleSpace for coordination state, scatter/gather, and shard groups, and worker pool. It supports native Rust as well a WASM runtime with SDKs for Python, Go, TypeScript, and Rust.
2. The Primitive Stack
PlexSpaces provides following set of coordination primitives:
| Primitive | What It Does |
|---|---|
| Shard Group | Partition data across N actors; scatter-gather with aggregation |
| Worker Pool | Stateless actor pool with load balancing |
| Process Group | Erlang pg2-style dynamic membership; broadcast |
| TupleSpace | Pattern-matched shared memory; Linda-model coordination |
| Channels | Queue-based stage coupling; 6 backends (Kafka, Redis, SQS, PG, …) |
| Workflow Actor | Multi-step durable orchestration; pause/resume/cancel |
| Distributed Lock | Lease-based mutual exclusion across actors |
| Ring AllReduce | Collective gradient reduction for distributed training |
| Parameter Server | Centralized gradient accumulation with pull/push |
| Broadcast | Send data to all actors in a process group |
| Collective Reduce | Sum/min/max across all actors; return to coordinator |
| Scatter/Gather | Fan-out to N workers, fan-in aggregated results |
2.1 Facets: Zero-Code Cross-Cutting Capabilities
PlexSpaces supports dynamic composition using Facets abstraction that allows plugging capabilities via configuration at deployment time:
| Facet | What It Adds |
|---|---|
security | JWT validation, RBAC enforcement |
logging | Structured log correlation with trace IDs |
metrics | Prometheus histograms per handler, per actor |
memoize | Deterministic response caching |
virtual_actor | Orleans-style activate-on-demand lifecycle |
schema_validation | Message schema enforcement |
durability | Journaling + checkpoint + crash replay |
execution_trace | Full execution trace capture |
timer / reminder | Durable delayed messages (survive crashes) |
event_sourcing / caching | Full audit trail / response caching |
kv / locks / registry / process_groups | State, coordination, discovery |
http_client | Outbound HTTP calls |
event_emitter | Pub/sub event broadcasting |

Behaviors:
| Behavior | Annotation | Pattern |
|---|---|---|
| GenServer | @actor / #[gen_server_actor] | Request-reply (ask/tell) |
| GenEvent | @event_actor | Fire-and-forget event handling |
| FSM | @fsm_actor | State machine with transitions |
| Workflow | @workflow_actor | Multi-step durable orchestration |
2.2 Polyglot SDK
PlexSpaces provides polyglot SDK for Python, TypeScript, GO and Rust:
Python:
from plexspaces import actor, handler, init_handler, host, state
@actor
class Counter:
count: int = state(default=0)
@init_handler
def on_init(self, config: dict):
self.count = int(config.get("args", {}).get("initial", 0))
@handler("increment")
def increment(self, amount: int = 1) -> dict:
self.count += amount
return {"count": self.count}
GO:
type CounterActor struct {
plexspaces.BaseActor
Count int `json:"count"`
}
func (a *CounterActor) Init(config map[string]interface{}) error {
if args, ok := config["args"].(map[string]interface{}); ok {
if v, ok := args["initial"].(float64); ok { a.Count = int(v) }
}
return nil
}
func (a *CounterActor) Handle(op string, payload map[string]interface{}) (interface{}, error) {
switch op {
case "increment":
amount := 1
if v, ok := payload["amount"].(float64); ok { amount = int(v) }
a.Count += amount
return map[string]interface{}{"count": a.Count}, nil
}
return nil, fmt.Errorf("unknown op: %s", op)
}
Typescript:
import { PlexSpacesActor, host } from "@plexspaces/sdk";
interface CounterState extends Record<string, unknown> {
count: number;
}
class CounterActor extends PlexSpacesActor<CounterState> {
onInit(config: Record<string, unknown>): CounterState {
const args = (config.args as Record<string, unknown>) || {};
return { count: Number(args.initial || 0) };
}
onIncrement(payload: Record<string, unknown>): Record<string, unknown> {
const amount = Number(payload.amount || 1);
this.state.count += amount;
return { count: this.state.count };
}
}
2.3 How Deployment Works
Every actor gets deployed via a TOML config file and a WASM binary. The config declares the supervision tree, actor roles, facets, and resource requirements:
[application]
name = "log-pipeline"
version = "1.0.0"
[supervisor]
strategy = "one_for_one"
max_restarts = 5
max_seconds = 60
[[supervisor.children]]
id = "LeaderActor"
actor_type = "LeaderActor"
behavior_kind = "GenServer"
role = "leader"
facets = ["virtual_actor", "metrics", "durability"]
[[supervisor.children]]
id = "PipelineWorker"
actor_type = "PipelineWorkerActor"
behavior_kind = "GenServer"
count = 8
role = "worker"
facets = ["virtual_actor", "metrics"]
The supervisor strategy (one_for_one, one_for_all, rest_for_one) comes directly from Erlang/OTP. If a pipeline worker crashes, only that worker restarts. If the leader crashes, the supervisor can optionally restart all children.
Deploy with a single API call:
zip app.zip pipeline_actor.wasm app-config.toml
curl -X POST http://node:8091/api/v1/applications/deploy \
-F "application_id=log-pipeline" \
-F "app_file=@app.zip"
3. Observability Pipelines
3.1 Log Ingestion & Routing Pipeline

In the log-pipeline example, the leader uses scatter/gather APIs to fan out event batches in parallel. Each PipelineWorkerActor receives a shard of events and runs them through a configurable function chain such as JSON parse, regex extract, PII masking, field enrichment, drop rules, etc. while keeping timing-metrics in state. The pipelines are partitioned by source hash across a shard group for horizontal scaling. After each scatterGather returns, the leader aggregates timing and throughput across all workers.
Python actor (excerpt):
@actor
class LeaderActor:
compute_ms: float = state(default=0.0)
coord_ms: float = state(default=0.0)
@handler("run")
def run(self, event_count=10000, worker_count=8, pipeline_fns=None, **kw):
pipeline_fns = pipeline_fns or ["json_parse", "regex_extract", "mask_pii", "enrich"]
group_id = host.create_shard_group(self.app_id, "PipelineWorkerActor", worker_count)
events = generate_log_events(event_count)
batches = partition_events(events, worker_count)
t0 = host.now_ms()
results = host.scatter_gather(group_id, "process_batch",
[{"events": b, "pipeline_fns": pipeline_fns} for b in batches], timeout_ms=60000)
coord_ms = host.now_ms() - t0
total_compute = sum(r.get("compute_ms", 0) for r in results)
return {"status": "ok", "events_per_sec": event_count * 1000 / coord_ms, ...}
Pipeline functions implemented across all languages:
fn_json_parse: Parse raw JSON bodyfn_regex_extract: Extract fields via regex (timestamps, IPs)fn_mask_pii: Mask email, SSN, credit card patternsfn_enrich: Add geo, datacenter, environment metadatafn_rename_fields: Normalize field names across sourcesfn_drop: Filter events by rules (debug level, internal sources)
Benchmark: Go WASM, 100K events, depth=5, 2 nodes (tested on my MacPro laptop)
| Metric | Value |
|---|---|
| Throughput | 177,304 events/sec |
| Wall time | 564 ms |
| Compute time | 1,585 ms (73%) |
| Coordination time | 564 ms (26%) |
| Granularity ratio | 2.8× |
| Events processed | 86,384 / 100,000 (13,616 dropped by fn_drop) |
Benchmark: Python WASM
| Workers | Total Events | Events/sec | Efficiency |
|---|---|---|---|
| 2 | 200 | ~10,000 | 100% |
| 4 | 400 | ~15,000 | 150% |
| 8 | 800 | ~22,000 | 220% |
| 16 | 1,600 | ~30,000 | 294% |
Benchmark: TypeScript WASM
| Workers | Events/sec | Wall ms | Comp ms | Coord ms | Comp% | Gran | Speedup | Eff% |
|---|---|---|---|---|---|---|---|---|
| 2 | 2,250 | 445 | 276 | 169 | 62% | 1.6× | 1.00× | 100% |
| 4 | 1,120 | 893 | 558 | 336 | 62% | 1.7× | 0.50× | 25% |
| 8 | 677 | 1,478 | 1,116 | 362 | 76% | 3.1× | 0.30× | 8% |
| 16 | 350 | 2,869 | 2,426 | 443 | 85% | 5.5× | 0.16× | 2% |
Go actor
func (a *WorkerActor) handleProcessBatch(payload map[string]interface{}) (interface{}, error) {
events := payload["events"].([]interface{})
fns := payload["pipeline_fns"].([]interface{})
t0 := plexspaces.NowMs()
var processed []map[string]interface{}
for _, ev := range events {
event := ev.(map[string]interface{})
for _, fn := range fns {
event = applyPipelineFunction(fn.(string), event)
if event == nil { break }
}
if event != nil { processed = append(processed, event) }
}
computeMs := plexspaces.NowMs() - t0
return map[string]interface{}{
"processed_count": len(processed),
"compute_ms": computeMs,
}, nil
}
3.2 Metrics Aggregation Pipeline

In the metrics_aggregation example, the leader uses scatter/gather APIs to provision N aggregation shards and then route metric batches. The framework hashes each metric name to a shard and each AggregationShard maintains a tumbling window in its actor state. After scatterGather returns, the leader collects shardResponses, merges rollup results, and runs z-score anomaly detection. Anomalies trigger alert events delivered via actor APIs and the state stays in the actor.
Rust embedded actor:
#[gen_server_actor]
struct AggregatorShard {
shard_id: usize,
windows: HashMap<String, Vec<f64>>,
total_metrics: u64,
}
#[plexspaces_handlers(gen_server)]
impl AggregatorShard {
#[handler("ingest_batch", cast)]
async fn handle_ingest_batch(&mut self, _ctx: &ActorContext, msg: &Message) -> Result<(), BehaviorError> {
let batch: Vec<Metric> = serde_json::from_slice(&msg.payload)?;
for m in &batch {
self.windows.entry(m.name.clone()).or_default().push(m.value);
}
self.total_metrics += batch.len() as u64;
Ok(())
}
#[handler("get_aggregates")]
async fn handle_get_aggregates(&self, _ctx: &ActorContext, _msg: &Message) -> Result<Value, BehaviorError> {
// Compute count/sum/avg/min/max per metric name, return as JSON
}
}
Benchmark: Go WASM (100K metrics, 8 workers)
| Metric | Value |
|---|---|
| Throughput | 480,769 metrics/sec |
| Compute time | 208 ms (69%) |
| Coordination time | 91 ms (30%) |
| Granularity ratio | 2.30× |
Benchmark: Rust Embedded
| Metric | Value |
|---|---|
| Throughput | 86,372 metrics/sec |
| Compute time | 361 ms (62%) |
| Coordination time | 220 ms (38%) |
| Granularity ratio | 1.64× |
Benchmark: Python WASM
| Workers | Total | Agg/s | Efficiency |
|---|---|---|---|
| 2 | 2K | 444K | 100% |
| 4 | 4K | 451K | 101% |
| 8 | 8K | 901K | 202% |
| 16 | 16K | 1.44M | 323% |
3.3 Distributed Tracing Pipeline

In the tracing_pipeline example, the leader creates a shard group and then uses scatter/gather to send only the spans whose trace_id hashes to that shard. All spans for a given trace land on one worker and each WorkerActor holds partial traces in its actor state and then runs tail-based sampling. Workers accumulate compute_ms in state and the leader sums shardResponses[*].compute_ms to compute aggregate compute time.
Go actor (excerpt):
type Span struct {
TraceID string `json:"trace_id"`
SpanID string `json:"span_id"`
ParentSpanID string `json:"parent_span_id,omitempty"`
ServiceName string `json:"service_name"`
Operation string `json:"operation"`
StatusCode int `json:"status_code"`
DurationMs float64 `json:"duration_ms"`
StartTimeMs int64 `json:"start_time_ms"`
Tags map[string]string `json:"tags,omitempty"`
}
Benchmark: Strong Scaling
| Workers | Spans/sec | Wall ms | Comp ms | Coord ms | Gran | Speedup | Eff% |
|---|---|---|---|---|---|---|---|
| 2 | 242,744 | 295 | 295 | 161 | 1.8x | 1.00x | 100.0% |
| 4 | 193,371 | 373 | 373 | 207 | 1.8x | 0.79x | 39.5% |
| 8 | 218,424 | 330 | 330 | 126 | 2.6x | 0.89x | 22.3% |
| 16 | 221,215 | 321 | 321 | 119 | 2.7x | 0.92x | 11.5% |
Service graph output example:
{
"nodes": [
{"name": "api-gateway", "span_count": 1000, "error_count": 48, "avg_latency_ms": 145.2},
{"name": "user-service", "span_count": 800, "error_count": 39, "avg_latency_ms": 32.1},
{"name": "order-service", "span_count": 600, "error_count": 31, "avg_latency_ms": 45.8},
{"name": "payment-service", "span_count": 400, "error_count": 22, "avg_latency_ms": 89.3}
],
"edges": [
{"source": "api-gateway", "target": "user-service", "call_count": 800, "p99_latency_ms": 195.0},
{"source": "api-gateway", "target": "order-service", "call_count": 600, "p99_latency_ms": 210.5},
{"source": "order-service", "target": "payment-service", "call_count": 400, "p99_latency_ms": 380.2}
]
}
Benchmark: Python WASM
| Workers | Total Traces | Spans/s | Gran | Eff% |
|---|---|---|---|---|
| 2 | 500 | 105,752 | 1.4× | 100% |
| 4 | 1,000 | 104,779 | 1.5× | 99% |
| 8 | 2,000 | 213,520 | 2.6× | 202% |
| 16 | 4,000 | 432,417 | 4.4× | 409% |
Benchmark: TypeScript WASM
| Shards | Total Spans | Span/s | Wall ms | Comp ms | Coord ms | Gran | Eff% |
|---|---|---|---|---|---|---|---|
| 2 | 1,000 | 5,988 | 167 | 109 | 74 | 1.5× | 100.0% |
| 4 | 2,000 | 5,764 | 347 | 227 | 150 | 1.5× | 96.3% |
| 8 | 4,000 | 6,920 | 578 | 455 | 175 | 2.6× | 115.6% |
| 16 | 8,000 | 7,105 | 1,126 | 987 | 245 | 4.0× | 118.7% |
4. ML Pipelines
4.1 Batch Inference

The parallel_ai_inference example demonstrates shard-group-based batch inference with model caching per worker. A benchmark actor then tries three ways of distributing that work and measures speed. The three ways are: broadcast like MPI, scatter/gather and elastic pool.

The batch_image_classification example shows ViT-style image classification across shards. A leader actor splits a pile of images across N worker actors, runs the work in rounds, and merges the results.
Benchmark: Python WASM Parallel AI Inference
| Shards | Req/sec | Comp% | Gran | Speedup | Efficiency |
|---|---|---|---|---|---|
| 1 | 512 | 62% | 1.7× | 1.00× | 100% |
| 4 | ~800 | — | — | ~1.56× | ~39% |
| 8 | 1,453 | 91% | 10.5× | 2.84× | 70% |
| 32 | ~2,100 | — | — | ~4.1× | 16% |
4.2 Parameter Server

The ParameterServerActor (leader) in parameter_server example creates shard group to provision worker shards, then broadcasts the current weight matrix via scatter/gather API. Each WorkerActor stores its local data shard in actor state, computes a mini-batch gradient locally, and returns a gradient summary. The leader collects shardResponses, averages the gradients across all workers.
Benchmark: TypeScript WASM
| Metric | Value |
|---|---|
| param_count | 6,464 (input_dim=100 × hidden_dim=64 + 64) |
| Gradient ops | workers × 20 rounds |
| Compute time | ~2 ms (gradient averaging + weight update) |
| Coordination time | ~4,000 ms (scatter-gather round trips) |
| Granularity ratio | ~0.00× (coordination-dominated on single node) |
| Errors | 0 |
4.3 Federated Learning

The federated_learning example shows a federated learning simulation. An aggregator coordinates many client actors, each holding its own private dataset. Every round, the aggregator sends the current model out. Each client trains locally and returns gradients. The aggregator then averages them weighted by sample count.
Benchmark: Federated Learning
| Clients | Wall ms | Compute ms | Coord ms | Gran | Speedup | Efficiency | Accuracy |
|---|---|---|---|---|---|---|---|
| 2 | 340 | 209 | 131 | 1.6× | 1.00× | 100% | 93.7% |
| 4 | 682 | 429 | 253 | 1.7× | 0.50× | 25% | 93.6% |
| 8 | 1,138 | 858 | 279 | 3.1× | 0.30× | 7.5% | 93.6% |
| 16 | 2,106 | 1,771 | 334 | 5.3× | 0.16× | 2.0% | 93.6% |
4.4 LLM Serving with Request Batching

The llm_serving simulates LLM serving pipeline similar to vLLM. A router generates a mix of requests and sends each one to a model tier: small, medium, or large. Cheap tasks go to the small model and heavy ones to the large mode. The router groups requests into batches per tier and scatters each batch to model workers, which simulate inference. The router then totals up throughput, tokens per second, cost, etc.
Benchmark: Python WASM
| Replicas | Req/s (aggregate) | Tok/s (aggregate) | Wall ms | Speedup |
|---|---|---|---|---|
| 2 | ~37,000 | ~2.1M | 54 | 1.00× |
| 4 | ~39,000 | ~2.2M | 51 | 1.05× |
| 8 | ~43,000 | ~2.5M | 46 | 1.17× |
Benchmark: TypeScript WASM
| Shards | Total Req | Req/s | Tok/s | Wall ms | Comp ms | Coord ms | Gran | Eff% |
|---|---|---|---|---|---|---|---|---|
| 2 | 1,552 | 86,222 | 23,325,778 | 18 | 25 | 5 | 2.6× | 100.0% |
| 4 | 3,096 | 96,750 | 27,481,500 | 32 | 49 | 20 | 0.6× | 112.2% |
| 6 | 4,644 | 132,686 | 38,290,286 | 35 | 73 | 23 | 0.5× | 153.9% |
| 8 | 6,144 | 170,667 | 44,575,111 | 36 | 97 | 24 | 0.5× | 197.9% |
4.5 Feature Store

The feature_store simulates a feature store like Feast or Tecton. A leader generates features and ingests the features into worker shards, which keep a versioned key-value store. The leader then serves batched lookups, where each worker checks its cache first and falls back to its store. At the end the leader collects stats and reports lookups per second, cache hit rate, etc.
TypeScript actor (excerpt):
class WorkerActor extends PlexSpacesActor<WorkerState> {
onIngestBatch(payload: Record<string, unknown>): Record<string, unknown> {
const records = payload.records as FeatureRecord[];
for (const rec of records) {
const key = `${rec.entity_id}:${rec.feature_name}`;
const versions = this.state.features[key] || { versions: [] };
versions.versions.push({ value: rec.value, version: rec.version, timestamp: rec.timestamp });
if (versions.versions.length > 3) versions.versions.shift(); // Keep last 3
this.state.features[key] = versions;
}
return { status: "ok", ingested: records.length };
}
onLookup(payload: Record<string, unknown>): Record<string, unknown> {
const entityId = payload.entity_id as string;
const cacheKey = entityId;
// Check LRU cache first
if (this.state.cache[cacheKey]) {
this.state.cache_hits++;
return { ...this.state.cache[cacheKey], cache_hit: true };
}
this.state.cache_misses++;
// Assemble feature vector from versioned store
// ...
}
}
Benchmark: Python WASM
| Workers | Lookups/sec | Speedup | Efficiency |
|---|---|---|---|
| 2 | ~10,000 | 1.00× | 100% |
| 4 | ~13,000 | 1.30× | 65% |
| 8 | ~17,000 | 1.70× | 43% |
| 16 | ~21,000 | 1.95× | 24% |
Benchmark: TypeScript WASM
| Shards | Total Lookups | Lkup/s | Wall ms | Comp ms | Coord ms | Gran | Eff% |
|---|---|---|---|---|---|---|---|
| 2 | 1,600 | 3,670 | 436 | 112 | 380 | 0.1× | 100.0% |
| 4 | 3,200 | 3,778 | 847 | 217 | 793 | 0.1× | 102.9% |
| 8 | 6,400 | 7,298 | 877 | 428 | 823 | 0.1× | 198.9% |
| 16 | 12,800 | 6,598 | 1,940 | 886 | 1,885 | 0.0× | 179.8% |
5. Data Processing & ETL
5.1 Streaming Window Aggregation

The streaming_pipeline example simulates streaming pipeline for event processing like a log or telemetry pipeline. A leader runs a series of batch rounds. Each round it tells all worker shards to process a batch of events. Workers simulate filtering, enriching, and transforming events. The worker then write a stage summary to a shared tuple space and return counts, bytes processed, etc. The leader reads the stage summaries back and reports events per second.
Benchmark: TypeScript WASM Streaming Pipeline
| Shards | Total Events | Evt/s | Wall ms | Comp ms | Coord ms | Gran | Eff% |
|---|---|---|---|---|---|---|---|
| 2 | 9,600 | 1,600,000 | 6 | 1 | 5 | 0.2× | 100.0% |
| 4 | 19,200 | 2,742,857 | 7 | 2 | 6 | 0.2× | 171.4% |
| 8 | 38,400 | 4,800,000 | 8 | 0 | 8 | 0.0× | 300.0% |
| 16 | 76,800 | 5,485,714 | 14 | 8 | 13 | 0.1× | 342.9% |
5.2 Data Lake RAG Ingestion

The data_lake_rag and agentic_rag_pipeline exampless demonstrate document ingestion, chunking, embedding, vector store and retrieval-augmented generation. Workflow actors manage the multi-step pipeline with checkpointing, so a crash during embedding resumes from the last chunk.
Benchmark: Data Lake RAG
| Metric | Value |
|---|---|
| Workload | 100 queries × 500 chunks × 2048 bytes |
| Wall time | 1,429 ms |
| Compute time | 3,984 ms (75.6%) |
| Coordination time | 1,285 ms (24.4%) |
| Granularity | 3.10× |
5.3 CDC Pipeline

The cdc_pipeline example simulates PostgreSQL-style WAL events (INSERT, UPDATE, DELETE) on four tables: users, orders, products, inventory. Workers maintain LSN (Log Sequence Number) tracking per table, apply schema transformations and fan out to three sink types (search index, analytics warehouse, cache invalidation).
Benchmark: CDC Pipeline (10K WAL events, 4 tables)
| Workers | Events/sec | Speedup | Efficiency |
|---|---|---|---|
| 2 | ~1,400 | 1.00× | 100% |
| 4 | ~1,900 | 1.36× | 68% |
| 8 | ~2,800 | 2.00× | 50% |
| 16 | ~3,100 | 2.20× | 28% |
Python CDC actor (excerpt):
@actor
class WorkerActor:
lsn_watermarks: dict = state(default_factory=dict)
seen_event_ids: set = state(default_factory=set)
@handler("process_events")
def process_events(self, events=None, **kw):
results = {"search_index": [], "analytics": [], "cache_invalidation": []}
for event in events:
if event["event_id"] in self.seen_event_ids:
continue # Dedup
self.seen_event_ids.add(event["event_id"])
self.lsn_watermarks[event["table"]] = max(
self.lsn_watermarks.get(event["table"], 0), event["lsn"])
transformed = self.transform(event)
for sink in ["search_index", "analytics", "cache_invalidation"]:
results[sink].append(self.format_for_sink(transformed, sink))
return {"status": "ok", "processed": len(events), "sinks": results}
5.4 Web Crawler

The web_crawl example simulates parallel web crawler. It includes aan orchestrator that runs a breadth-first crawl from seed URLs. For each URL it checks a fetcher using an elastic pool, fetches the page and checks the fetcher back in. After the crawl, the results are split between two analyzer shards that merge word counts into a global top-10. Also, the orchestrator asks the process group of fetchers how many pages each one handled.
Benchmark: TypeScript WASM Web Crawl
| Workers | TotPages | Pg/s | Wall ms | Comp ms | Coord ms | Speedup | Eff% |
|---|---|---|---|---|---|---|---|
| 2 | 400 | 1,951 | 205 | 200 | 5 | 1.00× | 50.0% |
| 4 | 800 | 1,923 | 416 | 202 | 214 | 0.99× | 24.6% |
| 8 | 1,600 | 3,883 | 412 | 201 | 211 | 1.99× | 24.9% |
| 16 | 3,200 | 7,692 | 416 | 201 | 215 | 3.94× | 24.6% |
6. Learnings
This post covered several use cases and examples for building data pipelines in different languages using PlexSpaces framework. Here are key lessons when building these pipelines:
- Stateful Actors. Actors are natural fit for building data pipelines. Facets allow adding cross-cutting concerns dynamically like metrics, durability, etc.
- Granularity ratio matters. There is a communication overhead when communicating with lots of actors/processes. The granularity ratio is compute time divided by coordination/communication time. So pay attention to this and test proper batch size for the parallel work
- Batch messages. Instead of sending one event per tell/cast, batch 100-1000 events per message.
- Broadcast scatter-gather. The scatter/gather APIs have been used in high performance computing for decades and are proven to scale for building highly parallel and concurrent applications.
- Reduce shard count when coordination dominates. More shards means more scatter-gather rounds. If your per-shard work is small, you are paying coordination tax for parallelism you do not need.
- Use cast (fire-and-forget) instead of call (request-reply). PlexSpaces supports both cast/tell and call/ask. Cast does not wait for a response, so the sender can continue immediately. Call blocks until the handler returns so use cast where possible such as data ingestion.
- TupleSpace for coordination. TupleSpace is based on Linda memory model, which is designed for coordination, synchronization.
- Durable actors. Stateless functions retry the whole batch on failure. Actors with the durability facet journal state before applying it. This means that a crash replays from the last journal entry instead of start.
- WASM Gotchas. WASM has high serialization overhead (~3.5ms for store reinstantiation). PlexSpaces supports native Rust SDK for microsecond-latency workloads.
Amdahl’s Law in practice. The parallel fraction of your workload determines the theoretical maximum speedup. From our benchmarks:
| Workload | Parallel Fraction | Max Speedup (Amdahl) | Measured at 16 workers |
|---|---|---|---|
| Log pipeline (batch=500) | 94% | 16.7x | 8.8x |
| Metrics aggregation | 86% | 7.1x | 5.2x |
| Batch inference | 97% | 33.3x | 12.0x |
| Ring AllReduce | 92% | 12.5x | 8.4x |
| Federated learning | 84% | 6.3x | 0.16× |
GitHub: https://github.com/bhatti/PlexSpaces
Clone the repo. Run ./test.sh in any example directory.
Related reading
- Building an Agent Harness and Eval Pipeline with Durable Actors
- Write a Redis Clone with Virtual Actors
- Building PlexSpaces: Decades of Distributed Systems Distilled
- Building Polyglot and Serverless Applications with WebAssembly
- Building Mini-OpenClaw: Secure AI Agents with Actors, WASM, and Supervision
- Building a Self-Improving AI Agent with Durable Actors — MiniHermes
- 20 Production Patterns for Distributed AI Agents Using Actors and TupleSpaces
- Multi-agent coordination patterns: Five approaches and when to use them
- Self-Improving AI Agent (MiniHermes)
- Redis Clone with Virtual Actors
- Agent Harness and Eval Pipeline
- Migrating off Cloudflare Durable Objects