NAME

EV::Kafka - High-performance asynchronous Kafka/Redpanda client using EV

SYNOPSIS

use EV::Kafka;

my $kafka = EV::Kafka->new(
    brokers  => '127.0.0.1:9092',
    acks     => -1,
    on_error => sub { warn "kafka: @_" },
    on_message => sub {
        my ($topic, $partition, $offset, $key, $value, $headers) = @_;
        print "$topic:$partition @ $offset  $key = $value\n";
    },
);

# Producer
$kafka->connect(sub {
    $kafka->produce('my-topic', 'key', 'value', sub {
        my ($result, $err) = @_;
        my $off = $result->{topics}[0]{partitions}[0]{base_offset};
        print "produced at offset $off\n";
    });
});

# Consumer (manual assignment)
$kafka->assign([{ topic => 'my-topic', partition => 0, offset => 0 }]);
my $poll = EV::timer 0, 0.1, sub { $kafka->poll };

# Consumer group
$kafka->subscribe('my-topic',
    group_id  => 'my-group',
    on_assign => sub { ... },
    on_revoke => sub { ... },
);

EV::run;

DESCRIPTION

EV::Kafka is a high-performance asynchronous Kafka client that implements the Kafka binary protocol in XS with EV event loop integration. It targets Redpanda and Apache Kafka (protocol version 0.11+).

Two-layer architecture:

  • EV::Kafka::Conn (XS) -- single broker TCP connection with protocol encoding/decoding, correlation ID matching, pipelining, optional TLS and SASL (PLAIN, SCRAM-SHA-256/512) authentication.

  • EV::Kafka::Client (Perl) -- cluster management with metadata discovery, broker connection pooling, partition leader routing, producer with key-based partitioning, consumer with manual assignment or consumer groups.

Features:

  • Binary protocol implemented in pure XS (no librdkafka dependency)

  • Automatic request pipelining per broker connection

  • Metadata-driven partition leader routing

  • Producer: acks modes (-1/0/1), key-based partitioning (murmur2), headers, fire-and-forget (acks=0), idempotent producer with epoch-bump recovery, transactional / exactly-once stream processing

  • Consumer: manual partition assignment, offset tracking, poll-based message delivery; consumer groups with JoinGroup/SyncGroup/Heartbeat, sticky partition assignment, offset commit/fetch, automatic rebalancing, session-expiry recovery

  • Compression: lz4, gzip, zstd, snappy (each gated by build-time library detection)

  • TLS (OpenSSL) and SASL/PLAIN, SCRAM-SHA-256/512 (with full RFC 5802 server-signature verification)

  • Automatic reconnection at the connection layer; bootstrap-broker failover; periodic metadata refresh

CLIENT LIFETIME

An EV::Kafka::Client is a handle, like an EV watcher: the application must keep a reference to it for as long as it is in use. When the last reference is dropped the client tears itself down immediately and silently, exactly as if close() had been called. This means

EV::Kafka->new(brokers => ...)->connect(sub { ... });
EV::run;

does NOT work -- the client is destroyed at the end of the statement and the connect callback never fires. Keep the object:

my $kafka = EV::Kafka->new(brokers => ...);
$kafka->connect(sub { ... });
EV::run;

The callback contract, both for explicit close() and for teardown by garbage collection:

  • No operation callback is ever silently dropped. Once an operation accepts a callback it is invoked exactly once. If the client is closed first, pending callbacks are invoked with an error in their documented error position: produce callbacks get (undef, $err), flush/poll callbacks get ($err), connect and seek callbacks get (undef, $err). Requests already sent to a broker are failed by the connection layer; queued, unsent work is failed by close() itself.

  • Event handlers are not operation callbacks. on_message, on_assign, on_revoke and on_connect simply stop firing at close; they are not invoked with an error.

  • Public methods croak after close. Calling any operation method on a closed client dies with "EV::Kafka: client is closed". close() itself is idempotent.

  • close() does not send LeaveGroup; call unsubscribe() first for a clean consumer-group leave.

RECONNECTION

Client-managed broker connections reconnect automatically after a loss (bootstrap-broker failover at connect() time; auto_reconnect with a 1s delay on every pooled connection afterwards). Operations queued while a broker is down complete when it returns. A lost connection is reported once via the on_error handler ("connection to broker lost; reconnecting"); the reconnect attempts themselves are not reported -- with the default die-on-error handler, a broker outage is survivable. Retries use capped exponential backoff with jitter. flush() bounds the wait: see flush_timeout.

ANYEVENT INTEGRATION

AnyEvent has EV as one of its backends, so EV::Kafka can be used in AnyEvent applications seamlessly.

NO UTF-8 SUPPORT

This module handles all values as bytes. Encode your UTF-8 strings before passing them:

use Encode;

$kafka->produce($topic, $key, encode_utf8($val), sub { ... });

CLUSTER CLIENT METHODS

new(%options)

Create a new EV::Kafka client. Returns a blessed EV::Kafka::Client object. Unknown option names croak ("EV::Kafka: unknown option(s): ..."), so typos fail at construction instead of being silently ignored.

my $kafka = EV::Kafka->new(
    brokers  => '10.0.0.1:9092,10.0.0.2:9092',
    acks     => -1,
    on_error => sub { warn @_ },
);

Options:

brokers => 'Str'

Comma-separated list of bootstrap broker addresses (host:port). Default: 127.0.0.1:9092.

client_id => 'Str' (default 'ev-kafka')

Client identifier sent to brokers.

tls => Bool

Enable TLS encryption.

tls_ca_file => 'Str'

Path to CA certificate file for TLS verification.

tls_skip_verify => Bool

Skip TLS certificate verification.

sasl => \%opts

Enable SASL authentication. Supported mechanisms: PLAIN, SCRAM-SHA-256, SCRAM-SHA-512.

sasl => { mechanism => 'PLAIN', username => 'user', password => 'pass' }
acks => Int (default -1)

Producer acknowledgment mode. -1 = all in-sync replicas, 0 = no acknowledgment (fire-and-forget), 1 = leader only.

linger_ms => Int (default 5)

Time in milliseconds to accumulate records before flushing a batch. Lower values reduce latency; higher values improve throughput.

batch_size => Int (default 16384)

Maximum batch size in bytes before a batch is flushed immediately.

compression => 'Str'

Compression type for produce batches: 'lz4' (requires liblz4), 'gzip' (requires zlib), 'zstd' (requires libzstd), 'snappy' (requires libsnappy), or undef for none.

idempotent => Bool (default 0)

Enable idempotent producer. Calls InitProducerId on connect and sets producer_id/epoch/sequence in each RecordBatch for exactly-once delivery (broker-side deduplication). Only one batch per (topic, partition) is in-flight at a time when this is enabled, to prevent sequence-number aliasing on retry. Produce calls made before InitProducerId completes are queued, never sent non-idempotent. If InitProducerId fails after bounded retries, the failure is sticky: queued and later produce callbacks receive an error and begin_transaction dies -- the client never silently degrades to non-idempotent production. DUPLICATE_SEQUENCE (46) responses are treated as acks (the broker already has the batch), while OUT_OF_ORDER_SEQUENCE_NUMBER (45) and INVALID_PRODUCER_EPOCH (47) trigger a single epoch-bump re-initialization shared by all affected partitions, then the batch is retried.

transactional_id => 'Str'

Enable transactional producer. Implies idempotent. Required for begin_transaction/commit_transaction/abort_transaction and send_offsets_to_transaction (full EOS). The transactional InitProducerId request is sent only to the transaction coordinator (discovered via FindCoordinator), never to an arbitrary broker.

partitioner => $cb->($topic, $key, $num_partitions)

Custom partition selection function. Default: murmur2 hash of key, or round-robin for null keys.

on_error => $cb->($errstr)

Error callback. Default: die.

on_connect => $cb->()

Called once after initial metadata fetch completes.

on_message => $cb->($topic, $partition, $offset, $key, $value, $headers)

Message delivery callback for consumer operations.

fetch_max_wait_ms => Int (default 500)

Maximum time the broker waits to accumulate fetch_min_bytes of data before returning a fetch response.

fetch_max_bytes => Int (default 1048576)

Maximum bytes per fetch response.

fetch_min_bytes => Int (default 1)

Minimum bytes before the broker responds to a fetch.

metadata_refresh => Int (default 300)

Periodic metadata refresh interval in seconds. Set to 0 to disable. Refreshes happen in the background, so consumers and producers pick up leader changes without waiting for a request to fail first.

flush_timeout => N (default 30)

Seconds flush() waits for outstanding work before giving up and calling its callback with an error. Set to 0 to wait forever.

loop => $ev_loop

EV loop object to use. Default: EV::default_loop.

connect([$cb])

Connect to the cluster. Connects to the first available bootstrap broker, fetches cluster metadata, then fires $cb->($metadata). On bootstrap-broker failure the next address is tried; if all fail, $cb is invoked as $cb->(undef, $err) (without a callback, the on_error handler fires instead). If the client is closed before the connection completes, $cb is invoked as $cb->(undef, $err).

Calling connect() more than once is defined: callbacks registered while a connection is in progress are queued and each fires exactly once; calling connect() on an already-connected client fires the callback immediately with the cached metadata. If idempotent producer initialization was requested but fails, queued connect callbacks are invoked as $cb->(undef, $err) instead of with metadata.

$kafka->connect(sub {
    my $meta = shift;
    # $meta->{brokers}, $meta->{topics}
});

produce($topic, $key, $value, [\%opts,] [$cb])

Produce a message. Routes to the correct partition leader automatically. %opts accepts only partition and headers; any other key croaks ("EV::Kafka: produce: unknown option ..."), catching typos like partiton.

# with callback (acks=1 or acks=-1)
$kafka->produce('topic', 'key', 'value', sub {
    my ($result, $err) = @_;
});

# with headers
$kafka->produce('topic', 'key', 'value',
    { headers => { 'h1' => 'v1' } }, sub { ... });

# fire-and-forget (acks=0)
$kafka->produce('topic', 'key', 'value');

# explicit partition
$kafka->produce('topic', 'key', 'value',
    { partition => 3 }, sub { ... });

produce_many(\@messages, $cb)

Produce multiple messages with a single completion callback. Each message is an arrayref [$topic, $key, $value] or a hashref {topic, key, value}. $cb fires when all messages are acknowledged.

$kafka->produce_many([
    ['my-topic', 'k1', 'v1'],
    ['my-topic', 'k2', 'v2'],
], sub {
    my $errors = shift;
    warn "some failed: @$errors" if $errors;
});

flush([$cb])

Flush all accumulated produce batches and wait for all in-flight requests to complete. $cb fires when all pending responses have been received, as $cb->($err) on failure: "flush timed out" after flush_timeout seconds (default 30, 0 waits forever), or "client closed" if the client is closed while a flush is outstanding.

assign(\@partitions)

Manually assign partitions for consuming.

$kafka->assign([
    { topic => 'my-topic', partition => 0, offset => 0 },
    { topic => 'my-topic', partition => 1, offset => 100 },
]);

seek($topic, $partition, $offset, [$cb])

Seek a partition to a specific offset. Use -2 for earliest, -1 for latest. Updates the assignment in-place. Errors are delivered as the callback's second argument: $cb->(undef, $err) fires when the partition is not assigned, when the partition leader is unreachable, or when the underlying ListOffsets request fails.

$kafka->seek('my-topic', 0, -1, sub {
    my ($res, $err) = @_;
    warn "seek failed: $err" if $err;
});

A seek that lands while a fetch for the same partition is in flight is not clobbered by the in-flight fetch response.

offsets_for($topic, $cb)

Get earliest and latest offsets for all partitions of a topic. The callback is $cb->($offsets, $err); $err is the first ListOffsets failure encountered, if any.

$kafka->offsets_for('my-topic', sub {
    my $offsets = shift;
    # { 0 => { earliest => 0, latest => 42 }, 1 => ... }
});

lag($cb)

Get consumer lag for all assigned partitions. The callback is $cb->($lag, $err); $err is the first ListOffsets failure encountered, if any.

$kafka->lag(sub {
    my $lag = shift;
    # { "topic:0" => { current => 10, latest => 42, lag => 32 } }
});

error_name($code)

Convert a Kafka numeric error code to its name. Callable as a method or class function.

$kafka->error_name(3);                  # "UNKNOWN_TOPIC_OR_PARTITION"
EV::Kafka::Client::error_name(3);       # same

poll([$cb])

Fetch messages from assigned partitions. Calls on_message for each received record. $cb fires as $cb->($err) when all fetch responses have arrived; $err is undef on success.

Per-partition fetch errors are surfaced and, where possible, recovered:

  • OFFSET_OUT_OF_RANGE (1): the position is re-resolved via ListOffsets according to auto_offset_reset ('earliest' for manually assigned consumers), and the poll callback receives the error. An out-of-range partition therefore recovers instead of failing every subsequent poll with the same stale offset.

  • Other error codes: reported to the on_error handler once per partition per rebalance (a new assign() or group rejoin resets the budget); the poll callback is not invoked with these.

An errored partition's position is never advanced by its error response. A die inside on_message is caught, reported via warn, and does not prevent delivery of the remaining records in the batch.

my $timer = EV::timer 0, 0.1, sub { $kafka->poll };

subscribe(@topics, %opts)

Join a consumer group and subscribe to one or more topics. The list of topic names comes first, followed by option key/value pairs. The group protocol handles partition assignment automatically.

Topic names must match the Kafka charset ([A-Za-z0-9._-]); anything else croaks. Option names are reserved words: group_id, group_instance_id, on_assign, on_revoke, session_timeout, rebalance_timeout, heartbeat_interval, auto_commit, auto_offset_reset. Because the API is positional, a misspelled option name cannot be detected (it would be a legal topic name), so avoid these strings as topic names. A reserved name given without a value croaks. group_id only needs to be defined, so group_id => '0' is accepted.

If the coordinator connection fails during JoinGroup/SyncGroup, the client retries with full coordinator re-discovery (5 attempts, 1s apart). On exhaustion the group enters the stopped state and the failure is reported via on_error instead of retrying forever.

$kafka->subscribe('topic-a', 'topic-b',
    group_id           => 'my-group',
    session_timeout    => 30000,      # ms
    rebalance_timeout  => 60000,      # ms
    heartbeat_interval => 3,          # seconds
    auto_commit        => 1,          # commit on unsubscribe (default)
    auto_offset_reset  => 'earliest', # or 'latest'
    group_instance_id  => 'pod-abc', # KIP-345 static membership
    on_assign => sub {
        my $partitions = shift;
        # [{topic, partition, offset}, ...]
    },
    on_revoke => sub {
        my $partitions = shift;
    },
);

commit([$cb])

Commit current consumer offsets to the group coordinator.

$kafka->commit(sub {
    my $err = shift;
    warn "commit failed: $err" if $err;
});

unsubscribe([$cb])

Leave the consumer group (sends LeaveGroup for fast rebalance), stop heartbeat and fetch loop. If auto_commit is enabled, commits offsets before leaving.

begin_transaction

Start a transaction. Requires transactional_id in constructor.

send_offsets_to_transaction($group_id, [$cb])

Commit consumer offsets within the current transaction via TxnOffsetCommit. This is the key step for exactly-once consume-process-produce pipelines.

$kafka->send_offsets_to_transaction('my-group', sub {
    my ($result, $err) = @_;
});

commit_transaction([$cb])

Commit the current transaction. All produced messages and offset commits within the transaction become visible atomically.

abort_transaction([$cb])

Abort the current transaction. All produced messages are discarded and offset commits are rolled back.

close([$cb])

Graceful shutdown: stop all timers, fail every pending operation callback with an error (see "CLIENT LIFETIME"), disconnect all broker connections. Idempotent. After close, calling any operation method croaks with "EV::Kafka: client is closed". Does not send LeaveGroup -- call unsubscribe() first for a clean consumer-group leave.

$kafka->close(sub { EV::break });

LOW-LEVEL CONNECTION METHODS

EV::Kafka::Conn provides direct access to a single broker connection. Useful for custom protocols, debugging, or when cluster-level routing is not needed.

my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
$conn->on_error(sub { warn @_ });
$conn->on_connect(sub { ... });
$conn->connect('127.0.0.1', 9092, 5.0);

connect($host, $port, [$timeout])

Connect to a broker. Timeout in seconds; 0 selects a 10-second default. The timeout covers the TCP connect AND every handshake phase (TLS, ApiVersions, SASL): a peer that accepts but then stays silent fails the connection with a "handshake timeout" error instead of hanging forever.

disconnect

Disconnect from broker.

connected

Returns true if the connection is ready (ApiVersions handshake complete).

metadata(\@topics, $cb)

Request cluster metadata. Pass undef for all topics.

$conn->metadata(['my-topic'], sub {
    my ($result, $err) = @_;
    # $result->{brokers}, $result->{topics}
});

produce($topic, $partition, $key, $value, [\%opts,] [$cb])

Produce a message to a specific partition.

$conn->produce('topic', 0, 'key', 'value', sub {
    my ($result, $err) = @_;
});

Options: acks (default 1), headers (hashref), timestamp (epoch ms, default now), compression ('none', 'lz4', 'gzip', 'zstd', 'snappy'; each requires its respective library at build time). With acks => 0 the broker sends no response, so the callback fires immediately with an empty success result: it means "handed to the socket", not "broker acknowledged".

produce_batch($topic, $partition, \@records, [\%opts,] [$cb])

Produce multiple records in a single RecordBatch. Each record is {key, value, headers}. Options: acks, compression, producer_id, producer_epoch, base_sequence.

$conn->produce_batch('topic', 0, [
    { key => 'k1', value => 'v1' },
    { key => 'k2', value => 'v2' },
], sub { my ($result, $err) = @_ });

fetch($topic, $partition, $offset, [\%opts,] $cb)

Fetch messages from a partition starting at $offset. %opts may set max_bytes (per-partition cap, default 1 MiB), max_wait_ms (broker block-time, default 500), min_bytes (default 1).

$conn->fetch('topic', 0, 0, sub {
    my ($result, $err) = @_;
    for my $rec (@{ $result->{topics}[0]{partitions}[0]{records} }) {
        printf "offset=%d key=%s value=%s\n",
            $rec->{offset}, $rec->{key}, $rec->{value};
    }
});

fetch_multi(\%topics, [\%opts,] $cb)

Multi-partition fetch in a single request. Groups multiple topic-partitions into one Fetch call to the broker. %opts accepts the same keys as fetch.

$conn->fetch_multi({
    'topic-a' => [{ partition => 0, offset => 10 },
                   { partition => 1, offset => 20 }],
    'topic-b' => [{ partition => 0, offset => 0 }],
}, sub { my ($result, $err) = @_ });

Used internally by poll() to batch fetches by broker leader, with max_bytes/max_wait_ms/min_bytes taken from the cluster client config.

list_offsets($topic, $partition, $timestamp, $cb)

Get offsets by timestamp. Use -2 for earliest, -1 for latest.

find_coordinator($key, $cb, [$key_type])

Find the coordinator broker. $key_type: 0=group (default), 1=transaction.

join_group($group_id, $member_id, \@topics, $cb, [$session_timeout_ms, $rebalance_timeout_ms, $group_instance_id])

Join a consumer group. Pass $group_instance_id for KIP-345 static membership.

sync_group($group_id, $generation_id, $member_id, \@assignments, $cb, [$group_instance_id])

Synchronize group state after join.

heartbeat($group_id, $generation_id, $member_id, $cb, [$group_instance_id])

Send heartbeat to group coordinator.

offset_commit($group_id, $generation_id, $member_id, \@offsets, $cb)

Commit consumer offsets.

offset_fetch($group_id, \@topics, $cb)

Fetch committed offsets for a consumer group.

api_versions

Returns a hashref of supported API keys to max versions, or undef if not yet negotiated.

my $vers = $conn->api_versions;
# { 0 => 7, 1 => 11, 3 => 8, ... }

on_error([$cb]), on_connect([$cb]), on_disconnect([$cb])

Set connection-level handler callbacks. Call with no argument or undef to clear.

If no on_error handler is installed, connection errors are reported via warn() rather than by throwing an exception from inside the event loop. on_disconnect fires exactly once per actual disconnection, so calling disconnect() from inside on_error or on_disconnect is safe and does not retrigger it. When a connection drops with requests still in flight, their callbacks fire with an error string BEFORE on_disconnect fires.

client_id($id)

Set the client identifier.

tls($enable, [$ca_file, $skip_verify])

Configure TLS.

sasl($mechanism, [$username, $password])

Configure SASL authentication.

auto_reconnect($enable, [$delay_ms])

Enable automatic reconnection. $delay_ms (default 1000) is the base delay: the first retry happens after exactly $delay_ms, later retries use capped exponential backoff (doubling each attempt, up to 30s, with +/-25% jitter), resetting once a connection succeeds. Reconnects reuse the timeout from the original connect() call.

leave_group($group_id, $member_id, $cb)

Send LeaveGroup to coordinator for fast partition rebalance.

create_topics(\@topics, $timeout_ms, $cb)

Create topics. Each element: {name, num_partitions, replication_factor}.

$conn->create_topics(
    [{ name => 'new-topic', num_partitions => 3, replication_factor => 1 }],
    5000, sub { my ($res, $err) = @_ }
);

delete_topics(\@topic_names, $timeout_ms, $cb)

Delete topics by name.

init_producer_id($transactional_id, $txn_timeout_ms, $cb)

Initialize a producer ID for idempotent/transactional produce. Pass undef for non-transactional idempotent producer.

add_partitions_to_txn($txn_id, $producer_id, $epoch, \@topics, $cb)

Register partitions with the transaction coordinator.

end_txn($txn_id, $producer_id, $epoch, $committed, $cb)

Commit ($committed=1) or abort ($committed=0) a transaction.

txn_offset_commit($txn_id, $group_id, $producer_id, $epoch, $generation, $member_id, \@offsets, $cb)

Commit consumer offsets within a transaction (API 28).

pending

Number of requests awaiting broker response.

state

Connection state as integer (0=disconnected, 6=ready).

Object lifetime

The connection is torn down when its last reference drops, or when DESTROY is called explicitly; explicit destruction is idempotent. Pending request callbacks are invoked with a 'destroyed' error during teardown. Once destroyed, the object is inert: every subsequent method call on it (from any copy of the reference) croaks with EV::Kafka::Conn: method called on destroyed connection -- including calls made from callbacks running during the teardown itself (those exceptions are caught by the callback dispatcher and reported via warn()). A conn passed a custom EV::Loop holds a reference on it, so the loop cannot be destroyed while the conn is alive.

UTILITY FUNCTIONS

EV::Kafka::_murmur2($key)

Kafka-compatible murmur2 hash. Returns a non-negative 31-bit integer.

EV::Kafka::_crc32c($data)

CRC32C checksum (Castagnoli). Used internally for RecordBatch integrity.

EV::Kafka::_error_name($code)

Convert Kafka error code to string name.

RESULT STRUCTURES

Produce result

$result = {
    topics => [{
        topic      => 'name',
        partitions => [{
            partition   => 0,
            error_code  => 0,
            base_offset => 42,
        }],
    }],
};

Fetch result

$result = {
    topics => [{
        topic      => 'name',
        partitions => [{
            partition      => 0,
            error_code     => 0,
            high_watermark => 100,
            records => [{
                offset    => 42,
                timestamp => 1712345678000,
                key       => 'key',      # or undef
                value     => 'value',     # or undef
                headers   => { h => 'v' },  # if present
            }],
        }],
    }],
};

Metadata result

$result = {
    controller_id => 0,
    brokers => [{ node_id => 0, host => '10.0.0.1', port => 9092 }],
    topics  => [{
        name       => 'topic',
        error_code => 0,
        partitions => [{
            partition  => 0,
            leader     => 0,
            error_code => 0,
        }],
    }],
};

ERROR HANDLING

Errors are delivered through two channels:

Connection-level errors fire the on_error callback (or croak if none set). These include connection refused, DNS failure, TLS errors, SASL auth failure, and protocol violations.
Request-level errors are delivered as the second argument to the request callback: $cb->($result, $error). If $error is defined, $result may be undef.

An error that belongs to a specific operation is reported to that operation's callback or to the on_error handler, never both -- this matters because the default on_error handler dies. Errors with no owning operation (connection loss, metadata failures, group protocol failures, per-partition fetch errors) go to on_error.

A die inside a user callback (on_message, a produce callback) does not propagate out of the event loop: it is caught, reported via warn, and the remaining callbacks in the same batch still run.

Within result structures, per-partition error_code fields use Kafka numeric codes:

0   No error
1   OFFSET_OUT_OF_RANGE
3   UNKNOWN_TOPIC_OR_PARTITION
6   NOT_LEADER_OR_FOLLOWER       (retried by the producer)
15  COORDINATOR_NOT_AVAILABLE    (retried)
16  NOT_COORDINATOR              (retried)
22  ILLEGAL_GENERATION           (group rejoin)
25  UNKNOWN_MEMBER_ID            (group rejoin)
27  REBALANCE_IN_PROGRESS        (group rejoin)
36  TOPIC_ALREADY_EXISTS
45  OUT_OF_ORDER_SEQUENCE_NUMBER (idempotent: epoch bump)
46  DUPLICATE_SEQUENCE_NUMBER    (idempotent: treated as an ack)
47  INVALID_PRODUCER_EPOCH       (idempotent: epoch bump)
79  MEMBER_ID_REQUIRED           (group rejoin with assigned id)

Use EV::Kafka::Client::error_name($code) for the full list.

When a broker disconnects mid-flight, all pending callbacks receive (undef, "connection closed by broker") or (undef, "disconnected").

ENVIRONMENT VARIABLES

These are used by tests and examples (not by the module itself):

TEST_KAFKA_BROKER    broker address for tests (host:port)
KAFKA_BROKER         broker address for examples
KAFKA_HOST           broker hostname for low-level examples
KAFKA_PORT           broker port for low-level examples
KAFKA_TOPIC          topic name for examples
KAFKA_GROUP_ID       consumer group for examples
KAFKA_LIMIT          message limit for consume example
KAFKA_COUNT          message count for fire-and-forget
BENCH_BROKER         broker for benchmarks
BENCH_MESSAGES       message count for benchmarks
BENCH_VALUE_SIZE     value size in bytes for benchmarks
BENCH_TOPIC          topic name for benchmarks

QUICK START

Minimal producer + consumer lifecycle:

use EV;
use EV::Kafka;

my $kafka = EV::Kafka->new(
    brokers    => '127.0.0.1:9092',
    acks       => 1,
    on_error   => sub { warn "kafka: @_\n" },
    on_message => sub {
        my ($topic, $part, $offset, $key, $value) = @_;
        print "got: $key=$value\n";
    },
);

$kafka->connect(sub {
    # produce
    $kafka->produce('test', 'k1', 'hello', sub {
        print "produced\n";

        # consume from the beginning
        $kafka->assign([{topic=>'test', partition=>0, offset=>0}]);
        $kafka->seek('test', 0, -2, sub {
            my $t = EV::timer 0, 0.1, sub { $kafka->poll };
            $kafka->{cfg}{_t} = $t;
        });
    });
});

EV::run;

COOKBOOK

Produce JSON with headers

use JSON::PP;
my $json = JSON::PP->new->utf8;

$kafka->produce('events', 'user-42',
    $json->encode({ action => 'click', page => '/home' }),
    { headers => { 'content-type' => 'application/json' } },
    sub { ... }
);

Consume from latest offset only

$kafka->subscribe('live-feed',
    group_id          => 'realtime',
    auto_offset_reset => 'latest',
    on_assign         => sub { print "ready\n" },
);

Graceful shutdown

$SIG{INT} = sub {
    $kafka->commit(sub {
        $kafka->unsubscribe(sub {
            $kafka->close(sub { EV::break });
        });
    });
};

At-least-once processing

$kafka->subscribe('jobs',
    group_id    => 'workers',
    auto_commit => 0,
);

# in on_message: process, then commit
on_message => sub {
    process($_[4]);
    $kafka->commit if ++$count % 100 == 0;
},

Batch produce

$kafka->produce_many([
    ['events', 'k1', 'v1'],
    ['events', 'k2', 'v2'],
    ['events', 'k3', 'v3'],
], sub {
    my $errs = shift;
    print $errs ? "some failed\n" : "all done\n";
});

Exactly-once stream processing (EOS)

my $kafka = EV::Kafka->new(
    brokers          => '...',
    transactional_id => 'my-eos-app',
    acks             => -1,
    on_message => sub {
        my ($t, $p, $off, $key, $value) = @_;
        my $result = process($value);
        $kafka->produce('output-topic', $key, $result);
    },
);

# consume-process-produce loop:
$kafka->begin_transaction;
$kafka->poll(sub {
    $kafka->send_offsets_to_transaction('my-group', sub {
        $kafka->commit_transaction(sub {
            $kafka->begin_transaction;  # next transaction
        });
    });
});

Topic administration

my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
$conn->on_connect(sub {
    $conn->create_topics(
        [{ name => 'new-topic', num_partitions => 6, replication_factor => 3 }],
        10000, sub { ... }
    );
});

BENCHMARKS

Measured on Linux with TCP loopback to Redpanda, 100-byte values, Perl 5.40.2, 20K messages (bench/benchmark.pl):

Pipeline produce (acks=1)   100K msg/sec    11.0 MB/s
Fire-and-forget (acks=0)    120K msg/sec    13.2 MB/s
Sequential round-trip        24K msg/sec    42 us avg latency
Metadata request             21K req/sec    47 us avg latency

Throughput by value size (pipelined, acks=1):

   10 bytes   105K msg/sec      1.5 MB/s
  100 bytes   100K msg/sec     10.5 MB/s
 1000 bytes    70K msg/sec     70.7 MB/s
10000 bytes    20K msg/sec    202.0 MB/s

Latency histogram (20K round-trips, acks=1, bench/latency.pl):

median: 39 us    p90: 59 us    p95: 75 us    p99: 122 us

Pipeline produce throughput is limited by Perl callback overhead per message. Fire-and-forget mode (acks=0) skips the response cycle entirely, reaching ~120K msg/sec. Sequential round-trip (one produce, wait for ack, repeat) measures raw broker latency around 39us median.

The fetch path is sequential (fetch, process, fetch again) which introduces one round-trip per batch. With larger max_bytes and dense topics, fetch throughput increases proportionally.

Run perl bench/benchmark.pl for throughput results. Set BENCH_BROKER, BENCH_MESSAGES, BENCH_VALUE_SIZE, and BENCH_TOPIC to customize.

Run perl bench/latency.pl for a latency histogram with percentiles (min, avg, median, p90, p95, p99, max).

KAFKA PROTOCOL

This module implements the Kafka binary protocol directly in XS. All integers are big-endian. Requests use a 4-byte size prefix followed by a header (API key, version, correlation ID, client ID) and a version-specific body.

Responses are matched to requests by correlation ID. The broker guarantees FIFO ordering per connection, so the response queue is a simple FIFO.

RecordBatch encoding (magic=2) is used for produce. CRC32C covers the batch from attributes through the last record. Records use ZigZag-encoded varints for lengths and deltas.

The connection handshake sends ApiVersions (v0) on connect to discover supported protocol versions. SASL authentication uses SaslHandshake (v1) + SaslAuthenticate (v2) with PLAIN mechanism.

Consumer group protocol uses sticky partition assignment with MEMBER_ID_REQUIRED (error 79) retry per KIP-394.

Non-flexible API versions are used throughout (capped below the flexible-version threshold for each API) to avoid the compact encoding complexity.

LIMITATIONS

  • Blocking DNS for hostnames -- numeric IPv4/IPv6 literals take a fast path (AI_NUMERICHOST) and never block. Non-literal hostnames call getaddrinfo synchronously, blocking the EV loop until the resolver responds. For fully non-blocking operation against named brokers, pre-resolve in Perl-land.

  • No GSSAPI/OAUTHBEARER -- only SASL/PLAIN and SCRAM-SHA-256/512 are implemented.

  • No flexible API versions -- all API versions are capped below the flexible-version threshold to avoid compact string/array encoding. Works with Kafka 0.11+ and Redpanda; loses access to a few newer protocol features.

  • Producer retry policy -- transient errors (NOT_LEADER, COORDINATOR_NOT_AVAILABLE) trigger metadata refresh and up to 3 retries with backoff. DUPLICATE_SEQUENCE is treated as an ack; OUT_OF_ORDER_SEQUENCE and INVALID_PRODUCER_EPOCH trigger one InitProducerId-with-fresh-epoch recovery attempt. Other broker errors are surfaced to the callback immediately.

CAVEATS

  • SIGPIPE on Linux with TLS -- writes on plain connections are protected (MSG_NOSIGNAL on Linux, SO_NOSIGPIPE on BSD/macOS), but TLS writes on Linux are not: SSL_write takes no per-call flags and Linux has no SO_NOSIGPIPE, so a broker RST racing a write can deliver SIGPIPE and terminate the process. Applications using tls => 1 on Linux should set $SIG{PIPE} = 'IGNORE'; at startup. Ignoring it is harmless: failed writes surface as EPIPE through the normal connection error path. The module deliberately does not set this itself -- a library must not change a process-global signal disposition its embedding application may rely on.

AUTHOR

vividsnow

LICENSE

This library is free software; you can redistribute it and/or modify it under the same terms as Perl itself.