Building a Data Transformation Plugin Pipeline
This page answers one task: a product processes streams of records — events, log lines, orders, sensor readings — and customers want to insert their own logic: filter, enrich, redact, reshape. You want a pipeline where each step is a WebAssembly plugin, so customers can upload code safely, with limits per step, clear failure handling, and throughput that keeps up with the stream.
Prerequisites
- [ ] A host runtime for plugins (Wasmtime with the Component Model, Extism, or similar) embedded in your service.
- [ ] A record format shared by all stages (JSON, MessagePack, or a typed WIT record).
- [ ] A queue or stream the pipeline reads from and writes to.
Why Wasm suits record pipelines
Record transformations are small, pure functions applied millions of times: take a record, return zero, one or several records. That fits WebAssembly’s strengths. Plugins run in-process with no network hop per record; they start in microseconds, so instances can be created per worker or even per batch; they are sandboxed, so customer code cannot read other customers’ data or the host’s memory; and limits on CPU (fuel or epochs) and memory are enforced by the runtime. Compared with running customer code in containers or separate processes, a Wasm pipeline has far lower per-record overhead, and compared with embedding a scripting language, it supports any language that compiles to Wasm.
The design questions are about the interface (what a stage sees and returns), batching (how many records per call), composition (how stages chain), and failure isolation (what happens when one stage traps or misbehaves).
Step 1 — define a batch-oriented interface
Define the stage interface in WIT so plugins in any language get typed bindings:
package acme:pipeline@1.0.0;
interface transform {
record record { key: string, payload: list, headers: list> }
variant outcome { keep(list), drop, fail(string) }
process: func(batch: list) -> list;
}
world stage {
import log: func(level: u8, msg: string);
export transform;
}
Passing a batch per call amortises the boundary cost — copying records into the component and results out — over many records. Returning one outcome per input
lets a stage keep, expand, drop or fail each record individually, so one bad record does not fail the whole batch.
Step 2 — run stages with limits
Instantiate each stage’s component per worker thread, reuse the instance across batches, and give every call a budget:
let mut store = Store::new(&engine, StageState::new(tenant_id));
store.limiter(|s| &mut s.limits); // StoreLimits: e.g. 64 MB memory cap per instance
store.set_fuel(FUEL_PER_BATCH)?; // CPU budget per call
let outcomes = match stage.acme_pipeline_transform().call_process(&mut store, &batch) {
Ok(o) => o,
Err(trap) => { metrics.trap(stage_id); return fail_batch(batch, trap); } // isolate and recreate instance
};
If a call traps or runs out of fuel, discard the instance and create a fresh one from the compiled component (microseconds), then route the batch to the dead-letter path. A misbehaving plugin affects only its own tenant’s records.
Step 3 — compose stages
Run stages in sequence for each batch: the outputs of stage n become the input batch of stage n+1. Keep records in the host’s memory between stages;
each stage copies in only what it needs. For long pipelines, consider composing several stages into one component with wac so records cross fewer
boundaries — at the cost of less granular isolation and limits. Configuration per tenant chooses which plugins form their pipeline and in which order.
Step 4 — handle failures explicitly
Define what happens on each failure type: a fail outcome for a single record sends it to a dead-letter queue with the stage name and message; a trap or fuel
exhaustion fails the batch for that stage, which can be retried record by record to isolate the poison record; repeated failures beyond a threshold disable the
stage for that tenant and alert them. Never silently drop records — every record ends in the output, the dead-letter queue, or an explicit drop outcome
that is counted.
Step 5 — observe each stage
Record per stage and tenant: records in and out, drops, failures, traps, fuel used per batch, latency per batch, and instance memory. These metrics show which plugin slows the pipeline, which tenant’s plugin keeps trapping, and whether a new plugin version changed behaviour. Expose a per-tenant view so customers can debug their own stages, including sample dead-letter records (with their consent and data rules).
Hot-swapping plugin versions
Customers update their plugins while the pipeline runs. Load the new version, run it against a sample of recent records (or the conformance suite) in a shadow mode, compare outputs with the current version, then switch at a batch boundary: new batches use the new instance, in-flight batches finish on the old one. Keep the previous version ready for instant rollback. Record which version processed each batch, so output differences can be traced to a deployment.
Throughput tuning
Batch size is the main lever: larger batches reduce per-call overhead but increase latency and memory per call. Measure records per second per core for batch sizes from 16 to 1,024 records and choose the knee of the curve for your latency target. Pre-serialise records in the format plugins use to avoid converting per stage, and pool instances per worker thread to avoid instantiation on the hot path.
Exactly-once and ordering guarantees
Pipelines sit between a source and a sink, and their delivery guarantees come from how stage results are committed, not from the plugins. Process a batch through all stages, write the outputs and dead letters, and only then acknowledge the input batch to the source — so a crash mid-batch causes reprocessing rather than loss. Reprocessing means plugins may see the same records twice; document that stages must be deterministic and free of external side effects, or give side-effecting stages idempotency keys derived from record keys. If downstream consumers rely on per-key ordering, partition batches by key and keep each partition on one worker, because parallel workers processing batches for the same key can reorder results. Plugins cannot fix any of this on their own; the host’s commit protocol defines the guarantees, and the plugin contract should state them plainly.
Per-tenant fairness
In a multi-tenant pipeline, one tenant’s slow plugin should not delay everyone else. Give each tenant its own queue or partition and schedule batches across tenants fairly — round-robin or weighted by plan — rather than processing a single shared queue in arrival order. Combine fuel limits per call with a per-tenant CPU budget per time window, and slow down (or pause) a tenant whose plugins consume more than their share, with clear metrics showing why. Fair scheduling plus strict per-call limits keeps the platform predictable even when individual plugins are expensive or buggy.
Expected output
The pipeline processes 180,000 records per second per core with three plugin stages at a batch size of 256; a tenant’s plugin that traps on a malformed record sends that record to the dead-letter queue while the rest of the batch proceeds after per-record retry; a plugin exceeding its fuel budget is disabled after repeated failures with an alert; and a new plugin version runs in shadow mode before taking over at a batch boundary.
Gotchas
- One record per call. Boundary overhead dominates. Batch.
- Failing whole batches for one bad record. Use per-record outcomes and retries.
- Reusing an instance after a trap. State is undefined. Recreate it.
- No dead-letter path. Records disappear. Account for every record.
- Swapping versions mid-batch. Inconsistent output. Switch at batch boundaries.
- One shared queue for all tenants. A slow plugin delays everyone. Schedule fairly per tenant.
Performance note
Increasing batch size from 1 to 256 records raised throughput from about 25,000 to about 180,000 records per second per core for a three-stage pipeline; beyond 512 the gain flattened while per-batch latency kept rising.
Frequently Asked Questions
Can stages run in parallel? Different batches can run in parallel on different workers; stages within one batch run in order.
Should stages be able to call external services? Only through host functions with allow-lists and timeouts; prefer pure transforms for throughput.
How do I let customers test plugins? Provide the conformance suite and a sandbox endpoint that runs their plugin on sample records.
Is JSON a good record format? It is easy for authors; binary formats are faster. Many pipelines pass bytes and let plugins choose.
Can plugins see the same record twice? Yes, after crashes and retries; stages should be deterministic and side-effect free, or use idempotency keys.
How do I keep one tenant from slowing others? Use per-tenant queues with fair scheduling, plus per-call fuel limits and a per-tenant CPU budget.
Related
- Testing plugins against a host contract — validating plugins.
- Limiting plugin CPU and memory use — limits.
- Hot reloading Wasm plugins — version swaps.
- Composing two Wasm components — composed stages.
← Back to Plugin Systems & Extensibility