Redpanda Explanation¶
What this page covers
How Redpanda works and why it is built this way: the Seastar thread-per-core reactor, Raft replication for every partition, the controller, Tiered Storage, Cloud Topics, Iceberg Topics, Wasm data transforms, Shadowing, the open-core licensing model, and the security model. Exact property defaults, versions, the Enterprise feature matrix and CVEs are in Reference. Step-by-step tasks are in How-to Guides.
Component Overview¶
Each broker is one C++ process. Seastar pins one reactor thread to each core, and each core (shard) owns a subset of partitions. Every partition is its own Raft group, and the controller is a special Raft group that holds cluster metadata. The diagram shows the components inside one broker and the object-storage tier it talks to.
flowchart TB
Client["Kafka clients / Redpanda Connect / kcat"]
subgraph Process["Redpanda broker process (one per node)"]
direction TB
subgraph Shards["Seastar shards (one reactor thread per core)"]
S0["Shard 0<br/>reactor + partitions"]
S1["Shard 1"]
SN["Shard N"]
end
KafkaAPI["Kafka API listener :9092"]
Controller["Controller Raft group<br/>redpanda/controller/0"]
PartitionRaft["Partition Raft groups<br/>(spread across shards)"]
StorageEngine["Storage engine<br/>(segments + indexes per NTP)"]
Archival["Archival uploader<br/>(Tiered Storage)"]
CTUploader["Cloud Topics batcher<br/>(object-first writes)"]
Datalake["Iceberg translator<br/>(Parquet + catalog commits)"]
SchemaRegistry["Schema Registry :8081<br/>(stored in _schemas topic)"]
Pandaproxy["HTTP Proxy :8082"]
AdminAPI["Admin API :9644"]
Wasm["Wasm transform runtime"]
end
subgraph ObjectStore["Object storage"]
Bucket["S3 / GCS / Azure Blob (ADLS Gen2)"]
end
Catalog["Iceberg catalog<br/>(object-storage or REST)"]
Client --> KafkaAPI
KafkaAPI --> Shards
Shards --> PartitionRaft
S0 --> Controller
PartitionRaft --> StorageEngine
StorageEngine --> Archival
PartitionRaft --> CTUploader
Archival --> Bucket
CTUploader --> Bucket
PartitionRaft --> Datalake
Datalake --> Bucket
Datalake --> Catalog
PartitionRaft --> Wasm
SchemaRegistry --> PartitionRaft
Pandaproxy --> PartitionRaft
AdminAPI --> Controller
Components¶
| Component | Role |
|---|---|
| Seastar reactor | One event loop per CPU core. Futures and promises replace blocking calls and locks. |
| Shard (core) | Owns a subset of partitions. Work moves between shards over lock-free message queues. |
| Controller | Cluster metadata: topic configuration, ACLs, users, roles, broker membership, partition placement, feature flags and license. |
| Per-partition Raft | Each partition is its own Raft group with its own leader and followers. |
| Storage engine | Append-only segments and indexes per namespace-topic-partition (NTP). |
| Archival uploader | Uploads segments and partition manifests to object storage for Tiered Storage. |
| Cloud storage read path | Serves old offsets from a local cache filled from object storage. |
| Cloud Topics | Topic class that writes record data to object storage first and replicates only metadata through Raft (GA in 26.1). |
| Iceberg translator | Converts topic data to Parquet and commits Iceberg snapshots to a catalog (Iceberg Topics). |
| Schema Registry | Confluent-compatible API for Avro, Protobuf and JSON Schema. Schemas live in the _schemas topic. |
| HTTP Proxy (Pandaproxy) | HTTP gateway for produce, consume and metadata. |
| Wasm transform runtime | Runs user-supplied WebAssembly modules that read an input topic and write to output topics. |
| rpk | CLI bundled with the broker. Talks to the Admin API on :9644 and the Kafka API. |
Thread-per-core (Seastar)¶
Traditional brokers (Kafka, RabbitMQ) use thread pools shared across CPU cores, which leads to context switches and lock contention. Redpanda uses Seastar:
- Each CPU core has one reactor thread pinned to it.
- Tasks on a core never block. Disk and network I/O are asynchronous and scheduled by Seastar's own I/O scheduler, which
rpk iotunecalibrates for the disk. - Shards talk to each other through lock-free message-passing queues.
- Memory is shard-local. Seastar pre-allocates most of the machine's memory at startup and splits it per core, which is why swap should be off and why memory is sized per core.
This removes lock contention on hot paths and gives predictable tail latency. The cost is extra bookkeeping when a request touches partitions on other shards. The broker balances partitions across cores at runtime (core_balancing_continuous), and since 24.3 the core count of a node can also be reduced.
Per-partition Raft¶
The sequence below shows an acks=all produce to a partition with replication factor 3.
sequenceDiagram
participant P as Producer (acks=all)
participant L as Partition leader (broker A)
participant F1 as Follower (broker B)
participant F2 as Follower (broker C)
P->>L: Produce(records)
L->>L: Append to log and fsync
par Replicate
L->>F1: append_entries
L->>F2: append_entries
end
F1-->>L: ack after fsync
Note right of L: majority (2 of 3) reached, commit index advances
L-->>P: ProduceResponse (offset)
F2-->>L: ack (late, still counted for later batches)
Kafka uses KRaft for metadata and ISR-based replication for data. Redpanda uses Raft for both. The benefits:
- One well-understood consensus protocol, with no "elections vs ISR" duality.
- A write acknowledged with
acks=allis on disk on a majority of replicas. The topic propertywrite.cachingrelaxes this to in-memory replication for lower latency. - Membership changes and partition moves are Raft reconfigurations.
The trade-off: every Raft group has a fixed cost (heartbeats, vote tracking), so clusters with very many tiny partitions are not the best fit. topic_partitions_per_shard caps the density.
Cluster Metadata (Controller)¶
The controller is a single Raft group (redpanda/controller/0) that covers the whole cluster. Its leader applies metadata changes and the log is replicated to the other brokers. It holds:
- Topic configurations, ACLs, SASL users and roles.
- Partition assignment to brokers and shards.
- Cluster membership and broker registration.
- Cluster configuration, feature flags and the Enterprise license.
The diagram shows how rpk or the Admin API changes metadata through the controller leader.
flowchart LR
RPK["rpk / Admin API :9644"]
subgraph ControllerGroup["Controller Raft group"]
Leader["Controller leader"]
Rep1["Controller replica"]
Rep2["Controller replica"]
end
Broker1["Broker 1<br/>partition manager"]
Broker2["Broker 2<br/>partition manager"]
Broker3["Broker 3<br/>partition manager"]
RPK --> Leader
Leader --> Rep1
Leader --> Rep2
Leader -->|"topic and placement deltas"| Broker1
Leader -->|"topic and placement deltas"| Broker2
Leader -->|"topic and placement deltas"| Broker3
Tiered Storage¶
Tiered Storage keeps the hot tail of each partition on local NVMe and moves older segments to object storage. It is an Enterprise feature in self-managed deployments (it is enabled by cloud_storage_enabled).
flowchart LR
LocalDisk["Local NVMe segments"]
Uploader["Archival uploader"]
Bucket["S3 / GCS / Azure"]
Manifest["Partition manifest<br/>(in the bucket)"]
Reader["Consumer fetch (old offset)"]
Cache["Cloud storage cache<br/>(local disk)"]
LocalDisk -->|"segment closed or upload interval hit"| Uploader
Uploader --> Bucket
Uploader --> Manifest
Reader --> Cache
Cache -->|"cache miss"| Bucket
A segment is uploaded when it closes, or at the latest after cloud_storage_segment_max_upload_interval_sec. Local retention can then be much shorter than total retention. Consumers that read old offsets are served from a local cache that is filled from the bucket on a miss. The on-bucket layout (segments plus Redpanda manifests) is Redpanda-specific. Apache Kafka's KIP-405 remote storage plugins cannot read it.
Remote Read Replicas (Enterprise) use the same data: a separate Redpanda cluster mounts a topic read-only from the bucket, so analytics consumers can read without loading the production cluster.
Cloud Topics¶
Cloud Topics (beta in 25.3, GA in 26.1) invert the storage model. Record data goes to object storage first, and only small metadata (object references and offsets) is replicated through Raft across zones. Produced batches are gathered per shard and uploaded as objects when a size threshold (cloud_topics_produce_batching_size_threshold, 4 MiB in current source) or an upload interval is reached.
sequenceDiagram
participant P as Producer
participant B as Broker shard (partition leader)
participant OS as Object storage (S3 / GCS / ADLS)
participant R as Raft followers (other AZs)
P->>B: Produce(batch)
B->>B: Buffer batch with other partitions on the shard
B->>OS: PUT object (size or time threshold)
OS-->>B: 200 OK
B->>R: Replicate placeholder metadata only
R-->>B: ack (majority)
B-->>P: ProduceResponse
The trade-off is latency for cost: produce latency includes an object-store PUT, but cross-AZ replication traffic, usually the largest line item for three-zone Kafka, mostly disappears (vendor claim: more than 90% less). Redpanda positions Cloud Topics for observability streams, analytics pipelines and ML feature feeds. Latency-sensitive topics stay on local or tiered storage in the same cluster. That is the "adaptable engine" (R1) idea behind 26.1. The storage mode is chosen per topic (local, tiered, cloud), with default_redpanda_storage_mode as the cluster default.
WASM Data Transforms¶
Data transforms run WebAssembly modules inside the broker. A transform reads records after they are written to its input topic. It is not inline in the produce path. For each record it can emit zero or more records to one or more output topics. Transforms are free in Community Edition and have been GA since 24.1. Enable them with data_transforms_enabled.
flowchart LR
Producer["Producer"] --> Input["Input topic partition"]
Input --> Transform["Wasm transform<br/>(runs next to the partition leader)"]
Transform --> Out1["Output topic A"]
Transform --> DLQ["Output topic B<br/>(for example, a dead-letter topic)"]
Out1 --> Consumer["Consumer"]
SDKs: Go (TinyGo), Rust, JavaScript and TypeScript (rpk transform init --language ...). Build with rpk transform build and deploy with rpk transform deploy. Transforms suit stateless, per-record work: filtering, masking, format conversion and routing. Joins, windows and state still need Flink, Redpanda Connect or another processor.
Iceberg Integration¶
Iceberg Topics (beta in 24.3, GA in 25.1 on 2025-04-07, Enterprise) let the broker write a topic's data into an Apache Iceberg table as well as the Kafka log. No Kafka Connect sink is needed.
flowchart LR
Producer["Producer"] --> Topic["Topic with redpanda.iceberg.mode set"]
SR["Schema Registry<br/>(Avro / Protobuf / JSON Schema)"] --> Translator
Topic --> Translator["Iceberg translator<br/>(per partition)"]
Translator --> Parquet["Parquet data files<br/>(Tiered Storage bucket)"]
Translator --> Catalog["Iceberg catalog<br/>(object_storage or REST: Glue, Unity, Snowflake Open Catalog)"]
Engines["Trino / Spark / Snowflake / Databricks / DuckDB"] --> Catalog
Engines --> Parquet
How it works:
iceberg_enabledturns the feature on for the cluster, andredpanda.iceberg.modeturns it on per topic.key_valuemode writes the raw key and value as binary columns.value_schema_id_prefixdecodes records that carry a Schema Registry ID (Confluent wire format) into typed columns.value_schema_latestuses the latest schema for the subject.- Records that cannot be translated follow
redpanda.iceberg.invalid.record.action(for example, a dead-letter table). - Commits happen on an interval, so the table lags the topic by roughly
redpanda.iceberg.target.lag.ms. - Since 25.2 JSON Schema is supported, as are AWS Glue, Databricks Unity Catalog and Snowflake Open Catalog as REST catalogs.
This collapses the usual Kafka to Connect to Parquet pipeline into one broker-native flow. The Iceberg files are open format and portable.
Shadowing (Disaster Recovery)¶
Shadowing (25.3, Enterprise, enable_shadow_linking) is Redpanda's built-in active-passive disaster recovery. A Shadow Link is created on the DR ("shadow") cluster and pulls from the source cluster. It replicates topic data with byte-for-byte offset fidelity, plus topic configs, consumer-group offsets, ACLs and Schema Registry data. On failover (rpk shadow failover), shadow topics become writable and clients re-point to the DR cluster with their committed offsets still valid.
stateDiagram-v2
[*] --> Replicating: rpk shadow create (on shadow cluster)
Replicating --> Replicating: pull data, offsets, ACLs, schemas
Replicating --> FailingOver: rpk shadow failover
FailingOver --> Promoted: shadow topics become writable
Promoted --> [*]: clients re-point to DR cluster
Unlike MirrorMaker 2, offsets are identical on both sides, so consumers need no offset translation. Replication is asynchronous, so a regional failover can still lose the last few seconds of writes (RPO greater than 0).
Redpanda Connect¶
Redpanda Connect is the former Benthos project, acquired in May 2024. It is a Go stream processor configured declaratively in YAML (inputs, a Bloblang processing pipeline, outputs) with hundreds of connectors, CDC inputs (Postgres, MySQL, MongoDB, SQL Server), Iceberg output, and an mcp-server mode that exposes pipelines as MCP tools for AI agents. It runs as a standalone binary, through rpk connect, as managed pipelines in Redpanda Cloud, or since 26.2 as an Operator Pipeline resource. Community and Certified connectors are Apache-2.0. Enterprise connectors need a license. The original Benthos engine remains MIT at redpanda-data/benthos, and WarpStream maintains a community fork called Bento.
Licensing Model¶
Redpanda is open core with source-available licenses, not OSI open source:
- Core broker: Redpanda BSL 1.1. Anyone may run it in production. The only exclusion is offering it to third parties as a "Streaming or Queuing Service", meaning a commercial offering where outsiders cause topics to be created. Each release converts to Apache-2.0 four years after its release date.
- Enterprise features: Redpanda Community License (RCL) plus a license key. The code ships in the same binary. Enabling a gated setting (Tiered Storage, Iceberg, Shadowing, RBAC, OIDC, audit logging and others) starts the 30-day trial clock or needs a purchased key.
The practical effect: a self-managed Community cluster is a full Kafka-API broker with local storage only. Anything that uses object storage is a paid feature, and Tiered Storage is free in Apache Kafka (KIP-405). Weigh that in TCO comparisons.
Security Model¶
Redpanda implements the Kafka security model (SASL, ACLs, TLS) and adds RBAC, group-based access control, OIDC and audit logging as Enterprise features.
Authentication¶
Kafka listeners support SASL/SCRAM (default), SASL/PLAIN (only alongside SCRAM), and mTLS in Community Edition. OAUTHBEARER (OIDC) and GSSAPI (Kerberos) are Enterprise. HTTP APIs (Admin, Schema Registry, HTTP Proxy) use HTTP Basic or OIDC (Enterprise). Redpanda Console signs users in with OIDC and maps them into cluster RBAC. The full mechanism table is in Reference, and setup steps are in How-to Guides.
Authorization¶
- Kafka ACLs apply to topics, groups, transactional IDs, the cluster, and Schema Registry subjects, with operations such as Read, Write, Create, Describe, Alter and Delete. Resource patterns are
literalorprefixed. - RBAC (Enterprise): named roles group ACLs, and users or groups are assigned to roles (
rpk security role). - GBAC (26.1): roles or permissions map to groups from the central IdP, so per-user grants are not needed on every cluster.
- Superusers (
superuserscluster property) bypass ACLs. Keep them to break-glass accounts.
Encryption¶
- In transit: TLS on every listener (Kafka, Admin, Schema Registry, HTTP Proxy, internal RPC), with optional client-certificate enforcement.
tls_min_versiondefaults to TLS 1.2. Kafka clients connect directly to each broker's advertised address, so a generic L7 proxy cannot sit in front of the Kafka protocol unless it is Kafka-aware. Terminate TLS on the Redpanda listeners. - At rest: local disks rely on OS or cloud volume encryption (LUKS, EBS, PD). Object-storage data relies on bucket-side encryption (SSE-S3 or SSE-KMS configured as the bucket default). The source has no Redpanda-level KMS key property.
- FIPS:
-fipsimage builds and the nodefips_modesetting (Enterprise).
Audit Logging¶
With audit_enabled (Enterprise), Redpanda writes audit events (management operations, authentication and optionally produce/consume access) to the internal _redpanda.audit_log topic. Because the log is a topic, you can ship it to a SIEM with any Kafka consumer or Redpanda Connect.
Threat Model¶
| Threat | Mitigation |
|---|---|
| ACL bypass through superusers | Keep superusers to break-glass accounts. Rotate them and audit their use. |
| Misconfigured object-storage bucket | Bucket-default SSE-KMS with a CMK, block-public-access, object-level access logging. |
| Console session hijack | Short OIDC token TTLs, HTTPS-only cookies. |
| Supply chain via third-party connectors | Pin Redpanda Connect and connector versions. Review plugins. Isolate on the network. |
| Unauthenticated Admin API | Enable TLS and authentication on :9644. It can change cluster config and users. |
| MITM on inter-broker RPC | Enable RPC TLS on production clusters. |
| Malicious or buggy Wasm transform | Transforms run inside the broker process, sandboxed by the Wasm runtime but sharing CPU and memory. Restrict who can deploy them. |
| Schema Registry spoofing | Require auth on the Schema Registry API. Pin the SR endpoint in clients. |
| Duplicate or replayed produce | Idempotent producers and transactions are on by default (enable_idempotence, enable_transactions). |
| OIDC token theft | Short TTLs, audience pinning, refresh-token rotation. |
| Accidental topic deletion | delete_topic_enable: false (Enterprise). Superusers cannot override it. |
| Regional outage | Shadowing to a second region, or Tiered Storage and Remote Read Replicas for read-only recovery. |
Comparison Hooks¶
- vs Kafka: same wire protocol. Redpanda wins on ops simplicity (single binary, no JVM) and tail latency. Kafka wins on ecosystem breadth, a fully Apache-2.0 license, and free KIP-405 tiered storage.
- vs Pulsar: Pulsar separates compute from storage in BookKeeper. Redpanda's Cloud Topics reach a similar cost profile per topic while keeping the Kafka API.
- vs NATS: different protocol. NATS is lighter for request-reply but not Kafka-compatible.
- Cross-broker comparison: Streaming Brokers Comparison.