|
This version is still in development and is not considered stable yet. For the latest stable version, please use Korvet 0.19! |
Kafka Compatibility
Korvet implements a subset of the Kafka protocol for compatibility with existing clients and tools.
Supported APIs
| API | Status | Notes |
|---|---|---|
Produce |
✅ Supported |
Send messages to topics |
Fetch |
✅ Supported |
Read messages from topics |
Metadata |
✅ Supported |
Topic and partition information |
ApiVersions |
✅ Supported |
Protocol version negotiation |
Consumer Groups |
✅ Supported |
JoinGroup, SyncGroup, Heartbeat, LeaveGroup, OffsetCommit, OffsetFetch |
Transactions |
❌ Not planned |
Use Redis transactions instead |
Admin API |
✅ Supported |
CreateTopics, CreatePartitions, DeleteTopics, DescribeConfigs, AlterConfigs, IncrementalAlterConfigs, DescribeCluster, ListGroups, DescribeGroups, DeleteGroups |
Idempotent Producers |
✅ Supported |
InitProducerId API for idempotent producer support (non-transactional) |
Kafka Version Compatibility
Korvet uses the Apache Kafka client library version 3.9.2 and is compatible with Kafka clients from version 2.8.0 and later.
Client Compatibility
-
Minimum supported client version: 2.8.0
-
Recommended client version: 3.9.x
-
Kafka client library: 3.9.2
Kafka clients are backward compatible, so newer clients (3.x, 4.x) can connect to Korvet without issues.
Protocol Features
Korvet implements Kafka protocol features equivalent to Kafka 2.8.0+, including:
-
Produce API (v0-v11; record batches must use the modern magic v2 format — legacy magic v0/v1 batches are rejected with
INVALID_RECORD) -
Fetch API (v0-v12, see Supported Version Matrix)
-
Consumer Group Protocol (JoinGroup, SyncGroup, Heartbeat, LeaveGroup)
-
Offset Management (OffsetCommit, OffsetFetch)
-
Topic Administration (CreateTopics, CreatePartitions, DeleteTopics)
-
Metadata API
Supported Version Matrix
Each API’s advertised max version is pinned explicitly to what the broker actually implements, rather than tracking whatever the bundled kafka-clients library defines. Kafka clients negotiate versions through the ApiVersions response, so any compatible client transparently uses the advertised version; features gated on a newer version fail fast on the client with UnsupportedVersionException instead of misbehaving on the broker.
| API | Advertised versions | Cap rationale |
|---|---|---|
Produce |
v0-v11 |
Transactional produce is not supported (transactional batches require a transaction coordinator, see Limitations). |
Fetch |
v0-v12 |
v13+ switches to topic-ID-based fetching (KIP-516) with incremental fetch sessions and follower fetching, none of which Korvet implements: topics are identified by name, each Fetch request is served statelessly, and there are no follower replicas. |
Metadata |
v0-v12 |
Request-side topic IDs (v10+) are not resolved; clients only fetch by topic ID on the Fetch path, which is capped at v12. |
ApiVersions |
v0-v4 |
|
FindCoordinator |
v0-v2 |
v3+ restricts the coordinator host/port response fields. |
JoinGroup |
v0-v9 |
|
SyncGroup |
v0-v5 |
|
Heartbeat |
v0-v4 |
|
LeaveGroup |
v0-v2 |
v3+ responses carry per-member errors ( |
OffsetCommit |
v0-v9 |
|
OffsetFetch |
v0-v9 |
Committed leader epoch is always |
ListOffsets |
v0-v6 |
The v7+ timestamp sentinels ( |
CreateTopics |
v0-v7 |
|
CreatePartitions |
v0-v3 |
|
DeleteTopics |
v0-v6 |
|
DescribeConfigs |
v0-v4 |
Config synonyms and documentation are not returned. |
AlterConfigs |
v0-v2 |
|
IncrementalAlterConfigs |
v0-v1 |
|
ListGroups |
v0-v5 |
|
DescribeGroups |
v0-v5 |
|
DeleteGroups |
v0-v2 |
|
DescribeCluster |
v0-v1 |
The v1 |
InitProducerId |
v0-v2 |
v3+ adds the KIP-360 epoch-bump handshake (client-supplied |
SaslHandshake |
v0-v1 |
|
SaslAuthenticate |
v0-v2 |
|
ConsumerGroupHeartbeat |
v0 |
Only advertised when the KIP-848 consumer protocol is enabled; v1+ (e.g. regex subscriptions) is not implemented. |
ConsumerGroupDescribe |
v0 |
Only advertised when the KIP-848 consumer protocol is enabled. |
The matrix is enforced in code: ApiVersionsHandler advertises these exact ranges, and ApiVersionsHandlerTest pins both the advertised matrix and the bundled library’s per-API latest versions, so a kafka-clients upgrade cannot silently raise an advertised version without an explicit review.
Leader Epoch for Stale Cache Detection
Korvet maintains a cluster-wide registry epoch that helps Kafka clients detect when their cached metadata is stale after broker membership changes:
-
Registry epoch: A shared counter incremented whenever the live broker set changes (broker join, leave, or stale heartbeat expiry)
-
Advertised as leader epoch: Each partition’s leader epoch in Metadata responses reflects the current registry epoch
-
Not-leader fetch responses: Include
currentLeaderwith the leader ID and epoch (Fetch v12+) so misrouted clients refresh metadata immediately
Single-broker mode: Always uses epoch 0.
Multi-broker clusters: The epoch increments on any membership change:
-
A new broker joins the cluster
-
An existing broker unregisters
-
A stale broker heartbeat expires
This allows clients to detect when their metadata cache is outdated and refresh it, reducing the number of misrouted requests after broker churn.
This is stale-cache detection only. Full KIP-320 leader epoch fencing (validating client-supplied currentLeaderEpoch to reject reads from older epochs) is not yet implemented.
|
Compression
Korvet supports all Kafka compression types:
-
NONE: No compression (default)
-
GZIP: Good compression ratio, higher CPU usage
-
SNAPPY: Balanced compression and speed
-
LZ4: Fast compression, lower CPU usage
-
ZSTD: Best compression ratio, moderate CPU usage
How Compression Works
Korvet implements server-side compression:
-
Producer side: Kafka clients can send compressed or uncompressed batches. Korvet automatically decompresses incoming batches into individual records before storing them.
-
Consumer side: When consumers fetch messages, Korvet compresses the response based on the topic’s
compression.typeconfiguration (not the producer’s compression setting). -
At rest: The producer’s Kafka batch compression is not retained — each record is stored as its own Redis Stream entry. The storage backend may then apply its own configurable at-rest compression to each record’s
valuefield (korvet.storage.local.compression.codec, defaultnone), independently of the Kafkacompression.typeused on the wire.
Configuring Compression
Compression is configured per-topic using the compression.type setting:
# Set compression for a topic (requires Admin API support)
kafka-configs --bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name my-topic \
--alter \
--add-config compression.type=lz4
The default compression type is NONE.
Limitations
-
Replication factor: Always 1 (Redis provides persistence)
-
Transactions: Not supported
-
Exactly-once semantics: Not supported (at-least-once delivery). Streaming frameworks that implement exactly-once sinks with Kafka transactions (for example Flink’s
KafkaSinkwithDeliveryGuarantee.EXACTLY_ONCE) fail at startup; use their at-least-once mode instead (see Apache Flink Integration). -
Consumer group membership: Held in broker memory. Committed offsets are durable in Redis, so after a broker restart
ListGroupsandDescribeGroupsreport groups with committed offsets in theEmptystate (with no member details) until clients rejoin.
Kafka Streams Compatibility
Korvet supports the broker-side features Kafka Streams relies on for group coordination and state:
-
Consumer group protocol, including the
CooperativeStickyAssignorused by Kafka Streams for task assignment -
Automatic creation of internal topics (repartition and changelog topics) via the Admin API
-
Log compaction (
cleanup.policy=compact) for changelog topics -
Idempotent producers (
InitProducerId)
Kafka Streams applications configured with processing.guarantee=at_least_once work against Korvet, verified by an integration test that runs a stateful groupByKey().count() topology and asserts:
-
Auto-creation of both classes of internal topic — the
-repartitiontopic and the-changelogtopic, the latter withcleanup.policy=compact -
State-store restoration from the changelog after an instance restarts with an empty state directory
-
A second instance joining the same
application.idand taking over part of the task assignment through the cooperative rebalance protocol, with the first instance retaining a subset of its tasks
Applications using processing.guarantee=exactly_once_v2 (the Kafka Streams default in recent versions) are not supported, because Korvet has no transaction coordinator: AddPartitionsToTxn, AddOffsetsToTxn, EndTxn, and TxnOffsetCommit are not implemented. Such an application fails during startup rather than degrading: InitProducerId with a transactional id returns UNSUPPORTED_VERSION, and the Streams client shuts down with StreamsException: Error encountered trying to initialize transactions. This work is tracked in issue #36.
Client Configuration
Most Kafka client configurations work with Korvet. Some settings are ignored:
-
acks: Always treated asacks=1 -
replication.factor: Ignored (always 1)
Unsupported topic configs, however, are rejected rather than ignored: a CreateTopics, AlterConfigs, or IncrementalAlterConfigs request that includes an unsupported key such as min.insync.replicas fails with INVALID_CONFIG ("Unknown or unsupported topic config"). Remove such keys from topic-creation tooling before pointing it at Korvet — tools that retry topic creation (for example the OpenMessaging Benchmark topic creator) can otherwise appear to hang. The accepted keys are listed in Topic Management.
Testing Compatibility
You can test Korvet with your existing Kafka applications by simply changing the bootstrap.servers configuration to point to Korvet.