Shahzad Bhatti Welcome to my ramblings and rants!

October 9, 2026

Building Production Data Pipelines with SEDA and Actors Model

Filed under: Computing — admin @ 8:44 pm

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:

PrimitiveWhat It Does
Shard GroupPartition data across N actors; scatter-gather with aggregation
Worker PoolStateless actor pool with load balancing
Process GroupErlang pg2-style dynamic membership; broadcast
TupleSpacePattern-matched shared memory; Linda-model coordination
ChannelsQueue-based stage coupling; 6 backends (Kafka, Redis, SQS, PG, …)
Workflow ActorMulti-step durable orchestration; pause/resume/cancel
Distributed LockLease-based mutual exclusion across actors
Ring AllReduceCollective gradient reduction for distributed training
Parameter ServerCentralized gradient accumulation with pull/push
BroadcastSend data to all actors in a process group
Collective ReduceSum/min/max across all actors; return to coordinator
Scatter/GatherFan-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:

FacetWhat It Adds
securityJWT validation, RBAC enforcement
loggingStructured log correlation with trace IDs
metricsPrometheus histograms per handler, per actor
memoizeDeterministic response caching
virtual_actorOrleans-style activate-on-demand lifecycle
schema_validationMessage schema enforcement
durabilityJournaling + checkpoint + crash replay
execution_traceFull execution trace capture
timer / reminderDurable delayed messages (survive crashes)
event_sourcing / cachingFull audit trail / response caching
kv / locks / registry / process_groupsState, coordination, discovery
http_clientOutbound HTTP calls
event_emitterPub/sub event broadcasting

Behaviors:

BehaviorAnnotationPattern
GenServer@actor / #[gen_server_actor]Request-reply (ask/tell)
GenEvent@event_actorFire-and-forget event handling
FSM@fsm_actorState machine with transitions
Workflow@workflow_actorMulti-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 body
  • fn_regex_extract: Extract fields via regex (timestamps, IPs)
  • fn_mask_pii: Mask email, SSN, credit card patterns
  • fn_enrich: Add geo, datacenter, environment metadata
  • fn_rename_fields: Normalize field names across sources
  • fn_drop: Filter events by rules (debug level, internal sources)

Benchmark: Go WASM, 100K events, depth=5, 2 nodes (tested on my MacPro laptop)

MetricValue
Throughput177,304 events/sec
Wall time564 ms
Compute time1,585 ms (73%)
Coordination time564 ms (26%)
Granularity ratio2.8×
Events processed86,384 / 100,000 (13,616 dropped by fn_drop)

Benchmark: Python WASM

WorkersTotal EventsEvents/secEfficiency
2200~10,000100%
4400~15,000150%
8800~22,000220%
161,600~30,000294%

Benchmark: TypeScript WASM

WorkersEvents/secWall msComp msCoord msComp%GranSpeedupEff%
22,25044527616962%1.6×1.00×100%
41,12089355833662%1.7×0.50×25%
86771,4781,11636276%3.1×0.30×8%
163502,8692,42644385%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)

MetricValue
Throughput480,769 metrics/sec
Compute time208 ms (69%)
Coordination time91 ms (30%)
Granularity ratio2.30×

Benchmark: Rust Embedded

MetricValue
Throughput86,372 metrics/sec
Compute time361 ms (62%)
Coordination time220 ms (38%)
Granularity ratio1.64×

Benchmark: Python WASM

WorkersTotalAgg/sEfficiency
22K444K100%
44K451K101%
88K901K202%
1616K1.44M323%

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

WorkersSpans/secWall msComp msCoord msGranSpeedupEff%
2242,7442952951611.8x1.00x100.0%
4193,3713733732071.8x0.79x39.5%
8218,4243303301262.6x0.89x22.3%
16221,2153213211192.7x0.92x11.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

WorkersTotal TracesSpans/sGranEff%
2500105,7521.4×100%
41,000104,7791.5×99%
82,000213,5202.6×202%
164,000432,4174.4×409%

Benchmark: TypeScript WASM

ShardsTotal SpansSpan/sWall msComp msCoord msGranEff%
21,0005,988167109741.5×100.0%
42,0005,7643472271501.5×96.3%
84,0006,9205784551752.6×115.6%
168,0007,1051,1269872454.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

ShardsReq/secComp%GranSpeedupEfficiency
151262%1.7×1.00×100%
4~800——~1.56×~39%
81,45391%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

MetricValue
param_count6,464 (input_dim=100 × hidden_dim=64 + 64)
Gradient opsworkers × 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)
Errors0

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

ClientsWall msCompute msCoord msGranSpeedupEfficiencyAccuracy
23402091311.6×1.00×100%93.7%
46824292531.7×0.50×25%93.6%
81,1388582793.1×0.30×7.5%93.6%
162,1061,7713345.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

ReplicasReq/s (aggregate)Tok/s (aggregate)Wall msSpeedup
2~37,000~2.1M541.00×
4~39,000~2.2M511.05×
8~43,000~2.5M461.17×

Benchmark: TypeScript WASM

ShardsTotal ReqReq/sTok/sWall msComp msCoord msGranEff%
21,55286,22223,325,778182552.6×100.0%
43,09696,75027,481,5003249200.6×112.2%
64,644132,68638,290,2863573230.5×153.9%
86,144170,66744,575,1113697240.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

WorkersLookups/secSpeedupEfficiency
2~10,0001.00×100%
4~13,0001.30×65%
8~17,0001.70×43%
16~21,0001.95×24%

Benchmark: TypeScript WASM

ShardsTotal LookupsLkup/sWall msComp msCoord msGranEff%
21,6003,6704361123800.1×100.0%
43,2003,7788472177930.1×102.9%
86,4007,2988774288230.1×198.9%
1612,8006,5981,9408861,8850.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

ShardsTotal EventsEvt/sWall msComp msCoord msGranEff%
29,6001,600,0006150.2×100.0%
419,2002,742,8577260.2×171.4%
838,4004,800,0008080.0×300.0%
1676,8005,485,714148130.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

MetricValue
Workload100 queries × 500 chunks × 2048 bytes
Wall time1,429 ms
Compute time3,984 ms (75.6%)
Coordination time1,285 ms (24.4%)
Granularity3.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)

WorkersEvents/secSpeedupEfficiency
2~1,4001.00×100%
4~1,9001.36×68%
8~2,8002.00×50%
16~3,1002.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

WorkersTotPagesPg/sWall msComp msCoord msSpeedupEff%
24001,95120520051.00×50.0%
48001,9234162022140.99×24.6%
81,6003,8834122012111.99×24.9%
163,2007,6924162012153.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:

  1. Stateful Actors. Actors are natural fit for building data pipelines. Facets allow adding cross-cutting concerns dynamically like metrics, durability, etc.
  2. 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
  3. Batch messages. Instead of sending one event per tell/cast, batch 100-1000 events per message.
  4. 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.
  5. 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.
  6. 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.
  7. TupleSpace for coordination. TupleSpace is based on Linda memory model, which is designed for coordination, synchronization.
  8. 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.
  9. 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:

WorkloadParallel FractionMax Speedup (Amdahl)Measured at 16 workers
Log pipeline (batch=500)94%16.7x8.8x
Metrics aggregation86%7.1x5.2x
Batch inference97%33.3x12.0x
Ring AllReduce92%12.5x8.4x
Federated learning84%6.3x0.16×

GitHub: https://github.com/bhatti/PlexSpaces

Clone the repo. Run ./test.sh in any example directory.

Related reading


    Powered by WordPress