redpanda-data / redpanda-data/connect

nats_jetstream input: allow configurable batch size for Fetch() to improve throughput

Open
#4,160 0 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

ux
Dominant language
Go
Stars
8.8k
Forks
969
Avg merge
1d 13h
Merged PRs (30d)
64

Description

Who is this for and what problem do they have today?

Users consuming high-volume JetStream streams where throughput matters more than per-message latency.

The nats_jetstream input hardcodes Fetch(1), pulling one message per network roundtrip. In our use case it caps throughput at ~2 MiB/min per consumer regardless of max_ack_pending. The only workaround is a broker with high copies, wasting goroutines and connections on sequential single-message fetches when NATS natively supports batch pulls.

What are the success criteria?

A new optional field (e.g. batch_size, defaulting to 1) passed to the NATS client's Fetch(batch int):

input:
    nats_jetstream:
      urls: ["nats://server:4222"]
      stream: my-stream
      durable: my-consumer
      batch_size: 256
Why is solving this problem impactful?

The NATS Go client's Fetch(batch int) already supports this - the hardcoded 1 leaves significant throughput on the table and forces users into workarounds. This would bring nats_jetstream in line with how other inputs handle batching (e.g. kafka with fetch_buffer_cap).

Additional notes

We ran into this while archiving ~2.5M DLQ messages to S3. Love the simplicity of the YAML pipeline approach - would be great to not need the broker workaround for this. Happy to submit a PR if that would be helpful!

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Locate the nats_jetstream input implementation and its configuration, then trace the hardcoded Fetch(1) call. Add an optional batch_size setting with a default of 1 and pass it to Fetch; done means the YAML setting controls batch pulls while existing configurations retain single-message behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
go
Domain
stream-processing
Issue type
Feature
Difficulty
3/5
Estimated time
1-2 days
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
68/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.