← all field notes
essay / September 23, 2026•22 min read

We diffed our pipeline against Vector's source. Here's what we found.

gorustperformanceelasticsearchkafka

We built a Logstash replacement in Go. Vector does twice the work on the same workload.

The pipeline is simple to describe: read log events from Kafka, parse and enrich them, and bulk-index them into Elasticsearch. We replaced Logstash with a Go service to get lower memory use, simpler deployment, and code we understand. We got all three. Then we hit a throughput wall.

We fixed the obvious problem, which took us from about 50k events/sec to about 1 lakh (100k). That figure is the total across every pod of both services, not a per-pod number. On the same workload, Vector does about 2 lakh (200k). We did not want to guess at more tuning. So we went through Vector's source and its documented defaults, line by line, and diffed them against our own code.

This post is that diff. It is not a "rewrite it in Rust" post. Most of what we found has nothing to do with the language. Each finding is a design choice that a production Rust pipeline made on purpose and that we made by accident. We are fixing each one in Go. The gap is still open at the end of this post, and we say so there.

1. The setup

The first version was one service that did everything: consume from Kafka, parse each event, build a bulk request, and push it to Elasticsearch. It was easy to reason about. It also meant parsing and pushing shared one process, one set of goroutines, and one failure domain.

↳ v1: one service does everything

kafkaraw log topics
→
consume + parse + bulk pushone process
→
elasticsearch

~50k events/sec, with consumer lag. A slow ES response stalled parsing, and a parsing spike delayed pushes.

The symptom was consumer lag. The cause was coupling. When Elasticsearch slowed down, the push goroutines blocked. Then the parse goroutines blocked behind them. Then the consumer stopped fetching. Parse-heavy bursts caused the opposite problem: pushes waited for CPU that parsing held.

2. The first fix: decoupling

We split the service in two. A parse/bridge service consumes raw events, parses them, and produces the parsed result back to Kafka. An ES-push service consumes the parsed topic and only does bulk indexing. Kafka sits between them as a durable buffer. Each side can now scale and fail on its own.

↳ v2: parse and push decoupled through kafka

kafkaraw log topics
→
parse/bridge serviceconsume + parse + produce
→
kafkalogs.parsed
→
ES-push serviceconsume + bulk index
→
elasticsearch

~100k events/sec in total, summed over every pod of both services. Vector on the same workload: ~200k.

Throughput doubled. It was a real fix, and it is the right architecture. But 100k, summed over every pod of two services, is still half of what Vector does on the same workload. The architecture was no longer the problem.

3. Why still short: going to the source

The usual next step is to tune: raise worker counts, raise batch sizes, add replicas, and watch a dashboard. We had done some of that, and each change moved the number a little or not at all. That is a sign you are tuning constants inside a design that caps you.

So we changed method. Vector solves the same problem, it is open source, and its defaults are documented. We read its Kafka source, its Elasticsearch sink, its batching and request layers, and its event model. For every place where Vector made a choice, we found the matching place in our code and wrote down what we did instead. Five groups of findings came out of that: concurrency, batching, regex, JSON, and copies. Concurrency is the biggest one.

4. Concurrency: a fixed semaphore vs. adaptive concurrency

What we do

The ES-push service limits in-flight bulk requests with one global semaphore, ESWorkers, default 16. Every partition of every topic shares it. There is a second problem in how the token is held. The retry loop sits inside the acquire/release pair, so a batch keeps its token through every backoff sleep:

go
var esSem = make(chan struct{}, cfg.ESWorkers) // 16, shared by every topic and partition

func (p *Pusher) push(ctx context.Context, batch []Event) error {
    esSem <- struct{}{}
    defer func() { <-esSem }()

    for attempt := 0; ; attempt++ {
        err := p.sendOnce(ctx, batch)
        if err == nil || attempt == p.maxRetries {
            return err
        }
        time.Sleep(backoff(attempt)) // token still held: one of 16 slots does nothing
    }
}

When Elasticsearch has a bad minute and starts returning 429s, batches start retrying. Each retrying batch sleeps while it holds a slot, so the effective limit drops below 16. Holding the slot is not the bug on its own. Vector does the same thing, as shown below. The bug is that nothing else reacts: new batches keep arriving at the same fixed limit while the cluster is asking everyone to slow down.

The parse/bridge service has the same shape on the produce side. Every worker, for every topic, produces through a hardcoded pool of 8 sarama.SyncProducers per broker. Each call blocks until every in-sync replica acknowledges, because RequiredAcks = WaitForAll:

go
cfg := sarama.NewConfig()
cfg.Producer.RequiredAcks = sarama.WaitForAll
cfg.Producer.Return.Successes = true // required by SyncProducer

const producersPerBroker = 8 // shared by every worker and every topic

func (p *pool) send(msgs []*sarama.ProducerMessage) error {
    prod := <-p.idle          // wait for one of 8 producers
    defer func() { p.idle <- prod }()
    return prod.SendMessages(msgs) // blocks for the full ISR round trip
}

WaitForAll is the correct durability setting, and we are keeping it. The problem is the pool. Eight synchronous callers means at most eight produce round trips in flight per broker, whatever the worker count is. Everything else queues behind the pool.

What Vector does

Vector's HTTP-based sinks, the Elasticsearch sink included, default to Adaptive Request Concurrency (ARC). The shared sink service type shows the layering. Note that Retry sits inside the concurrency limit, the same as our retry loop inside the semaphore:

pub type Svc<S, L> =
    RateLimit<AdaptiveConcurrencyLimit<Retry<FibonacciRetryPolicy<L>, Timeout<S>>, L>>;

The difference is the limit itself. The defaults start at 1 and cap at 200:

const fn default_initial_concurrency() -> usize {
    1
}

const fn default_decrease_ratio() -> f64 {
    0.9
}

const fn default_ewma_alpha() -> f64 {
    0.4
}

const fn default_rtt_deviation_scale() -> f64 {
    2.5
}

const fn default_max_concurrency_limit() -> usize {
    200
}

Once per RTT window, the controller compares the window's mean RTT with an EWMA of past RTTs. It adds 1 when the limit was fully used and RTT was at or below the average. It multiplies by 0.9 when RTT goes past the average by 2.5 standard deviations, or when any response in the window was back-pressure:

if inner.current_limit < self.settings.max_concurrency_limit
    && inner.reached_limit
    && !inner.had_back_pressure
    && current_rtt.is_some()
    && current_rtt.unwrap() <= past_rtt.mean
{
    // Increase (additive) the current concurrency limit
    self.semaphore.add_permits(1);
    inner.current_limit += 1;
}
// Back pressure responses, either explicit or implicit due
// to increasing response times, trigger a decrease in the
// concurrency limit.
else if inner.current_limit > 1
    && (inner.had_back_pressure || current_rtt.unwrap_or(0.0) >= past_rtt.mean + threshold)
{
    let new_limit =
        ((inner.current_limit as f64 * self.settings.decrease_ratio) as usize).max(1);
    self.semaphore
        .forget_permits(inner.current_limit - new_limit);
    inner.current_limit = new_limit;
}

"Back-pressure" means any response the sink would retry. For Elasticsearch, that is a 429 or any 5xx:

match status {
    StatusCode::TOO_MANY_REQUESTS => RetryAction::Retry("too many requests".into()),
    StatusCode::NOT_IMPLEMENTED => {
        RetryAction::DontRetry("endpoint not implemented".into())
    }
    _ if status.is_server_error() => RetryAction::Retry(
        // ...
    ),
    // ...
}

This is AIMD, the same additive-increase/multiplicative-decrease loop that TCP congestion control uses (the TCP post has a visualizer for that one). The limit finds the cluster's real capacity, follows it down when the cluster slows, and climbs back when it recovers. Nobody picks the number.

On the Kafka side, Vector's sink uses rust-rdkafka's FutureProducer. send_result only enqueues the record into librdkafka's internal queue and returns a future for the delivery report. librdkafka batches and pipelines the actual produce requests. When the queue is full, the sink waits 100ms and tries again. No small pool of blocking callers exists to become the bottleneck:

loop {
    match this.kafka_producer.send_result(record) {
        // Record was successfully enqueued on the producer.
        Ok(fut) => {
            drop(blocked_state.take());
            return fut
                .await
                .expect("producer unexpectedly dropped")
                // ...
        }
        // Producer queue is full or a policy has been violated and the request should
        // be retried
        Err((
            KafkaError::MessageProduction(
                RDKafkaErrorCode::QueueFull | RDKafkaErrorCode::PolicyViolation,
            ),
            original_record,
        )) => {
            // ...
            record = original_record;
            tokio::time::sleep(Duration::from_millis(100)).await;
        }
        // ...
    }
}

See the difference

The simulator below runs both strategies against the same model cluster. Start it at 1× and watch ARC climb past 16 toward the cluster's real capacity while the fixed limit stays flat. Then move the slowdown slider up and watch the fixed limit become too high: RTT grows, 429s arrive, and retries hold slots. Then add transient 5xx errors. ARC counts those as back-pressure too, so it backs off hard, and at low slowdown the fixed limit can win. ARC is built to protect the cluster, not to win every benchmark.

↳ fixed ESWorkers=16 vs. adaptive request concurrency

—
fixed-16 req/s
—
ARC req/s
1
ARC in flight
64
cluster's real capacity

press start

1×
0%

cluster: 40ms base RTT, 64 concurrent bulk requests before it queues. The fixed limit of 16 leaves 48 slots of headroom unused — no amount of load will ever push past it.

ARC: ramping from concurrency 1

The lesson is not "16 is the wrong number." Every constant is the wrong number most of the time, because the cluster's capacity changes during the day. A fixed limit is too low when the cluster is healthy and too high when it is not. The fix in Go is to replace the constant with a feedback loop. An AIMD limiter is about 60 lines of Go. We are also moving the backoff sleep outside the limiter. Vector does not do that, so that part is our choice, not a finding from the diff. On the produce side, we are replacing the sync pool with sarama.AsyncProducer, and we are evaluating franz-go, which pipelines produce requests natively.

go
for attempt := 0; ; attempt++ {
    if err := p.limiter.Acquire(ctx); err != nil { // adaptive limit, not a fixed channel
        return err
    }
    start := time.Now()
    err := p.sendOnce(ctx, batch)
    p.limiter.Release(time.Since(start), isRetryable(err)) // feed RTT and back-pressure back
    if err == nil || attempt == p.maxRetries {
        return err
    }
    time.Sleep(backoff(attempt)) // no slot held while sleeping
}

5. Batching economics

Every bulk request pays a fixed cost before Elasticsearch indexes a single document: the network round trip, HTTP framing, and fan-out from the coordinating node to the shards. Batch size decides how many events share that cost.

Our ES-push service caps a batch at 250 events or 2MB, whichever comes first. With log events of a couple of KB, the 250-event cap always binds first, and the 2MB cap never matters. Vector's Elasticsearch sink uses RealtimeSizeBasedDefaultBatchSettings: 10MB, no event-count cap, and a 1-second timeout. It limits by bytes only, because bytes are what the cluster pays for.

/// Reasonable default batch settings for sinks with timeliness concerns, limited by byte size.
#[derive(Clone, Copy, Debug, Default)]
pub struct RealtimeSizeBasedDefaultBatchSettings;

impl SinkBatchSettings for RealtimeSizeBasedDefaultBatchSettings {
    const MAX_EVENTS: Option<usize> = None;
    const MAX_BYTES: Option<usize> = Some(10_000_000);
    const TIMEOUT_SECS: f64 = 1.0;
}
pub batch: BatchConfig<RealtimeSizeBasedDefaultBatchSettings>,

Small batches hurt more when concurrency is fixed. With 16 requests in flight and 250 events in each, at most 4,000 events are in flight at any time, however fast the cluster is. Move the slider to see where the fixed overhead stops dominating.

↳ batch economics — who pays the per-request overhead

93k
events/s
93%
of each request spent on fixed overhead
0.5MB
bulk body size
286k143k0kours: 250vector: 10MBevents/s vs. events per bulk request (16 in flight, 2KB events)
250

parse/bridge workers vs. partitions — routing is partition % workers

12

12 workers can ever receive work. 84 sit idle forever, each still waking on a 10ms flush timer — 8,400 empty wakeups/s that buy nothing.

The second panel shows a related problem in the parse/bridge service. The worker count is hardcoded to 96 or 256, whatever the topic's partition count is. Messages go to worker partition % workers. So a topic with 12 partitions uses exactly 12 workers. The other 84 never get a message. Each one still runs a 10ms flush ticker, wakes up, finds an empty buffer, and goes back to sleep.

go
workers := make([]*worker, max(cfg.Workers, 96)) // floor of 96, no link to partitions

func (r *router) dispatch(m *sarama.ConsumerMessage) {
    w := r.workers[int(m.Partition)%len(r.workers)] // only len(partitions) workers ever match
    w.in <- m
}

func (w *worker) run() {
    t := time.NewTicker(10 * time.Millisecond) // fires on idle workers too
    for {
        select {
        case m := <-w.in:
            w.add(m)
        case <-t.C:
            w.flush() // empty flush on most workers, forever
        }
    }
}

Real parallelism was capped at the partition count all along. The extra workers added only idle goroutines and timer wakeups. The fix is to size workers as min(partitions, GOMAXPROCS × k) and to read the partition count from the consumer group claim, not from config.

6. The regex detour: same algorithm, different constant factor

This one started as "why is Vector's field extraction faster on the same patterns?" We expected an algorithm difference. There is none. Go's regexp and Rust's regex crate are both RE2-style engines. Both guarantee linear time in the input, and neither backtracks exponentially. A pattern that is safe in one is safe in the other.

The gap comes from constant factors, and most of them were in our code:

  • Compiling per call. This is the most common Go regex mistake. A regexp.MustCompile inside a function that runs per event costs tens of microseconds and dozens of allocations, every event.
  • Capture groups when you only need a match. FindStringSubmatch allocates a result slice and forces Go to use a slower engine path than MatchString.
  • Unanchored .* prefixes. A leading .* removes the literal prefix that lets the engine skip ahead, so it scans every byte.
  • Engine-level differences. These are real, but smaller than the three above. Rust's crate has a lazy DFA and SIMD literal prefilters (memchr, Teddy). Go's engine uses an NFA, a one-pass matcher, or a bounded backtracker, with a simpler literal prefix check.
go
// before: compiled on every event
func parseLine(line string) (map[string]string, bool) {
    re := regexp.MustCompile(`.*status=(\d+) latency=(\d+)ms`)
    m := re.FindStringSubmatch(line)
    ...
}

// after: compiled once, anchored on a literal
var lineRE = regexp.MustCompile(`^\S+ \S+ status=(\d+) latency=(\d+)ms`)

func parseLine(line string) (map[string]string, bool) {
    m := lineRE.FindStringSubmatch(line)
    ...
}

Toggle each choice below. With compile-per-call turned on, Rust is slower than Go, because Regex::new does more up-front work so that matching is cheap later. The language is not the variable here.

↳ same regex, same algorithm class — where the constant factor goes

time per log line
Go
20,410 ns
Rust
45,381 ns
heap allocations per log line
Go
42
Rust
61

Compile-per-call dominates both sides — Rust's is actually the more expensive compile. At 1 lakh events/s this alone is 2.0 core-seconds per second in Go. Fix it first; nothing else on this card matters until you do.

7. Death by a thousand JSON round-trips

This is where the Go code spent most of its CPU. Each event passes through encoding/json several times, and most of those passes do no useful work.

The bulk body: buildBulk

For each event, buildBulk decodes the whole document into a map[string]json.RawMessage to read one or two fields. Then it calls json.Marshal, which uses reflection, on a nested map to produce the action line. That line is a fixed template with two variable values:

go
// before: reflection to build a fixed string
action, _ := json.Marshal(map[string]any{
    "index": map[string]any{"_index": index, "_id": id},
})

// after: index names are [a-z0-9._-], ids are generated hex — no escaping needed
buf.WriteString(`{"index":{"_index":"`)
buf.WriteString(index)
buf.WriteString(`","_id":"`)
buf.WriteString(id)
buf.WriteString("\"}}\n")

To be fair to Go: Vector does not use a string template here either. It builds the action line with serde_json's json! macro. That still allocates a small value tree, but it involves no runtime reflection, which is the expensive part of json.Marshal on a map:

(true, DocumentMetadata::Id(id)) => {
    write!(
        writer,
        "{}",
        json!({
            bulk_action: {
                "_index": index,
                "_id": id,
            }
        }),
    )
}

The bulk response: sendOnce

After each request, sendOnce decodes the entire bulk response into a map[string]json.RawMessage, then calls json.Unmarshal once for each item. That is 251 decode calls for a 250-item batch. It happens on every request, including the normal case where every item succeeded. Elasticsearch already reports this in a top-level errors boolean:

go
var head struct {
    Errors bool `json:"errors"` // items is skipped without being decoded
}
if err := json.Unmarshal(body, &head); err != nil {
    return err
}
if !head.Errors {
    return nil // fast path: the whole batch succeeded
}
return p.collectItemErrors(body) // per-item decode only when something failed

Vector's Elasticsearch retry logic goes further. On a 2xx response it does not decode anything unless a plain substring search finds a failure:

_ if status.is_success() => {
    let body = simdutf_bytes_utf8_lossy(response.http_response.body());

    if body.contains("\"errors\":true") {
        match EsResultResponse::parse(&body) {
            // ...
        }
    }
    // ...
}

The nested message: parseSaramaMessage

The worst one is in the parse/bridge service. Each Kafka record wraps the real log event as a JSON string in a "message" field. parseSaramaMessage decodes that string into any. That is the most allocation-heavy path in encoding/json: every object becomes a map, every array a slice, and every value is boxed in an interface. Then the function passes the result straight to json.Marshal and never reads it.

go
// before: build a full tree, then serialize it straight back
var inner any
if err := json.Unmarshal([]byte(env.Message), &inner); err != nil {
    return nil, err
}
out, err := json.Marshal(inner) // re-encodes what it just decoded

// after: check the syntax, keep the bytes
raw := json.RawMessage(env.Message)
if !json.Valid(raw) {
    return nil, errInvalidJSON
}

The two versions do not produce byte-identical output. json.Marshal on a map sorts keys and escapes <, > and &. The RawMessage version forwards the original bytes as they are. Both are equal as JSON. Before you switch, check that nothing downstream compares documents as bytes.

Step through the payload to see where the allocations come from:

↳ allocations per token — interface{} vs. struct vs. raw bytes

{"ts":"2026-09-01T10:00:00Z""level":"error""svc":"checkout""tags":["db""retry"]"http":{"status":504"ms":1203}}json.Marshal(v)

var v any + re-marshal

0
heap allocations

typed struct

0
heap allocations

json.Valid + RawMessage

scans the bytes to check syntax, builds nothing
0
heap allocations

Step through the nested "message" payload. Red squares are allocations the generic decode makes, one token at a time.

What Vector does

Vector's JSON decoder is not magic. It parses into a generic serde_json::Value tree, which is the Rust equivalent of our decode into any, and then converts that tree into an event:

let json: serde_json::Value = match self.lossy {
    true => serde_json::from_str(&simdutf_bytes_utf8_lossy(&bytes)),
    false => serde_json::from_slice(&bytes),
}
.map_err(|error| format!("Error parsing JSON: {error:?}"))?;

// If the root is an Array, split it into multiple events
let mut events = match json {
    serde_json::Value::Array(values) => values
        .into_iter()
        .map(|json| Event::from_json_value(json, log_namespace))
        .collect::<Result<SmallVec<[Event; 1]>, _>>()?,
    _ => smallvec![Event::from_json_value(json, log_namespace)?],
};

The difference is what happens next. Vector keeps that tree as the event for the rest of the pipeline, and the sink serializes it once, straight into the bulk body with serde_json::to_writer. Our pipeline builds a tree, serializes it, sends the bytes through Kafka, and decodes them again in the ES-push service. Vector does not win on parser speed. It wins because it does not repeat the round-trip.

8. Copies, copies, copies

On the way out of the parse/bridge service, each payload is copied up to three times:

go
payload, _ := json.Marshal(out)           // 1: allocates the payload (needed)
msg := cloneSaramaMessage(&sarama.ConsumerMessage{Value: payload})
                                           // 2: deep copy; payload is never used again
...
func WriteSaramaBatch(msgs []*sarama.ConsumerMessage) {
    for _, m := range msgs {
        v := make([]byte, len(m.Value))
        copy(v, m.Value)                   // 3: copies again for the producer message
        batch = append(batch, &sarama.ProducerMessage{Topic: topic, Value: sarama.ByteEncoder(v)})
    }
}

The first allocation is necessary. Copies 2 and 3 protect against aliasing that cannot happen, because nothing else holds a reference to payload. At 100k events/sec and 2KB each, those two copies alone write about 400MB/sec of memory that becomes garbage at once. That raises GC frequency, and GC takes CPU from parsing.

The Grow call that does not add up

buildBulk tries to size its buffer before it writes, which is the right idea. But it calls Grow once per event:

go
var buf bytes.Buffer
for _, ev := range events {
    buf.Grow(len(ev) + actionLineLen) // intent: reserve space for the whole batch
}
// actual result: capacity for about the single largest event

Grow(n) guarantees space for n more bytes past the current length. The loop writes nothing, so the length stays 0. Each call asks for room for one event, and the buffer ends up sized for the largest one. The write loop that follows then grows and copies the buffer several times, which is what the pre-sizing was meant to prevent. The fix is to add up the sizes first and call Grow once:

go
total := 0
for _, ev := range events {
    total += len(ev) + actionLineLen
}
buf.Grow(total)

What Vector does

Vector makes copies cheap in two places. Raw input arrives as bytes::Bytes, a reference-counted view into a shared buffer: the decoder above takes bytes: Bytes, and cloning a Bytes increments a counter without copying memory. The event itself is also reference-counted. A LogEvent holds its fields behind an Arc:

#[derive(Clone, Debug, Default, Deserialize, PartialEq)]
pub struct LogEvent {
    #[serde(flatten)]
    inner: Arc<Inner>,
    // ...
}

pub fn value_mut(&mut self) -> &mut Value {
    let result = Arc::make_mut(&mut self.inner);
    // ...
}

#[derive(Clone)] on a struct that holds an Arc means cloning an event bumps a reference count. The fields are copied only when someone writes to a shared event, through Arc::make_mut. This is copy-on-write: a fan-out to three sinks costs three counter increments, not three deep copies.

Go has no Bytes type in the standard library, but it does not need one here. Pass the slice and agree on who owns it. Go slices already share memory when you do not copy them.

↳ one payload, two pipelines — copies vs. a shared handle

go parse/bridge service

json.Marshal(event)…
cloneSaramaMessage(m)…
WriteSaramaBatch(msgs)…
producer.SendMessages…

buffers allocated: 0 × 2KB

vector pipeline

source: Bytes::from(buf)…
transform: event.clone()…
sink: encode + batch…
request sent, dropped…

buffers allocated: 0 × 2KB · refcount —

614 MB/s
Go: payload bytes allocated + copied
410 MB/s
of which pure waste (becomes garbage)
2KB
100k

9. What we're actually going to change

One line per finding, in the order we plan to ship them:

  1. Replace the fixed ESWorkers semaphore with an AIMD limiter that uses RTT and retryable responses as feedback.
  2. Release the limiter slot during retry backoff. This one is our own choice. Vector keeps retries inside its limit.
  3. Replace the pool of 8 sync producers with sarama.AsyncProducer, or move to franz-go.
  4. Drop the 250-event cap and batch by bytes only, starting at 5–10MB, with a timeout.
  5. Size parse/bridge workers from the partition count and CPU count, not a hardcoded floor.
  6. Hoist every regexp.MustCompile to package level, and anchor patterns on a literal.
  7. Build the bulk action line from a string template, not json.Marshal.
  8. Check the bulk response's top-level errors first, and decode items only when it is true.
  9. Pass the nested message as json.RawMessage after json.Valid, with no decode to any.
  10. Remove the two redundant payload copies, and fix the Grow loop to call Grow once with the total.

We have not re-benchmarked yet, so the gap is still open. We are at about 100k events/sec, summed over every pod of both services, and Vector is at about 200k on the same workload. We do not expect every item on this list to matter equally. The concurrency changes and the JSON changes are the ones we expect to move the number most. We will publish the measured result as a follow-up, whether it closes the gap or not.

The method is the main thing we took away. When tuning stops working, find a system that solves the same problem well, read its source, and write down every place where its design differs from yours. Most of what you find will not be about the language. It will be about constants that nobody chose on purpose.


Related: Kafka beyond the basics covers the partition and consumer group mechanics that decide how many workers can actually get work. Go channels: what's actually inside hchan covers semaphore sizing with chan struct{}, which is the limiter that section 4 replaces. TCP from the inside has the congestion-control loop that ARC borrows from.