kafkaex/kafka_ex
67.0
Adequate · 3 October 2026
18.2k
lines of production code
Elixir
primary language
2
measurements over time
What this system is
KafkaEx is an Elixir client library for interacting with Apache Kafka, providing robust APIs for producing and consuming messages, managing consumer groups, and administering topics. It supports modern Kafka features including SASL authentication (PLAIN, SCRAM, OAuthBearer, MSK IAM), SSL/TLS, and comprehensive telemetry for observability. The system ensures reliability through configurable retry logic, dynamic supervision, and stable consumer group coordination with rebalancing and heartbeat management.
How it got here
2014–2017 — v1.0 API overhaul and modernization
10 changes.
This period focused on the major v1.0 release, introducing a new dynamic supervision-based client API, comprehensive telemetry for observability, and modern Elixir configuration syntax. The work also included significant infrastructure improvements such as Credo linting, expanded test coverage with Mimic and Testcontainers, and utility scripts for local SSL/TLS development.
2023–2026 — KafkaEx v1.0 release and Kayrock migration
27 changes.
This period focused on the KafkaEx v1.0 release, which involved migrating the protocol layer to the Kayrock library and introducing a comprehensive new consumer group architecture with stream-based APIs. The work also included implementing robust SASL authentication mechanisms, refining cluster metadata handling, and establishing extensive chaos and integration testing to ensure reliability under network stress.
Features
Add Credo linter configuration and Elixir formatter settings
The repository now includes \.credo.exs\ and \.formatter.exs\ to enforce code quality and consistency. Credo is configured with a standard set of checks for consistency, readability, refactoring, and warnings, while the formatter is set to a 120-character line length and includes the \mix.exs\, \config\, \lib\, and \test\ directories.
(repo-wide) · high confidence
Added structs to parse consumer group member assignments
New modules have been added to the \ConsumerGroupDescription\ message family to represent the detailed assignment information for individual consumer group members. Specifically, \Member\, \MemberAssignment\, and \PartitionAssignment\ structs now model the data returned by the Kafka DescribeGroups API, including member identifiers, client metadata, and the specific topic-partition assignments. This enables users to inspect exactly which partitions are assigned to which consumers within a group.
_lib/kafka\_ex/messages/consumer\_group\description · high confidence
Introduce Kayrock-compatible Kafka client implementation
The \lib/kafka\_ex/client\ directory now contains a new \KafkaEx.Client\ GenServer that implements the legacy KafkaEx.Server API using the Kayrock protocol. This change adds a complete request/response pipeline—including \RequestBuilder\, \ResponseParser\, and \RequestContext\ modules—to handle Kafka operations such as produce, fetch, and consumer group coordination. The client introduces a structured retry budget system (\RequestBudget\) to align per-attempt socket timeouts with total call timeouts, and adds a \NodeSelector\ to support routing strategies like \:node\_id\, \:random\, and \:controller\. It also includes \MetadataLog\ for throttling missing-topic warnings and \Error\ structs for standardized error reporting, ensuring that the client remains compatible with existing KafkaEx API callers while leveraging Kayrock's protocol handling.
_lib/kafka\ex/client · high confidence
New consumer group and stream APIs for Kafka consumption
This change introduces a new consumer architecture in the \lib/kafka\_ex/consumer\ directory, providing \KafkaEx.Consumer.ConsumerGroup\ for broker-coordinated group membership and partition assignment, \KafkaEx.Consumer.GenConsumer\ as a behavior for implementing custom message handlers with configurable offset commit strategies (sync/async), and \KafkaEx.Consumer.Stream\ for functional stream-based consumption. It also adds \KafkaEx.Consumer.GenConsumer.Supervisor\ for managing static consumer processes. These modules replace or supplement the previous low-level API, offering higher-level abstractions for handling fetch retries, offset management, and group coordination.
_lib/kafka\ex/consumer · high confidence
New message structs for Kafka Admin and Consumer APIs
Added new message structs to support Kafka Admin and Consumer operations, including ApiVersions for broker capability negotiation, CreateTopics and DeleteTopics for topic management, and consumer group structures (ConsumerGroupDescription, ConsumerGroupListing, JoinGroup, SyncGroup, LeaveGroup, Heartbeat, FindCoordinator) for group coordination. Also added Fetch for record retrieval with control-batch handling, Offset/PartitionOffset/OffsetAndMetadata for offset management, RecordMetadata for produce acknowledgments, Header for message metadata, and Heartbeat for consumer liveness.
_lib/kafka\ex/messages · high confidence
New testing and development utility scripts for KafkaEx
Added a set of shell scripts to the \scripts/\ directory to support local testing and development workflows. \docker\_up.sh\ automates starting the Kafka environment using Docker Compose, with architecture-aware selection (ARM64 vs x86\_64). \gen\_ssl.sh\ generates self-signed SSL certificates and keystores (PKCS\#12) for local SSL/TLS testing. \check\_teardown.sh\ enforces safe process teardown practices in tests by detecting racy \GenServer.stop\ or \Supervisor.stop\ calls in \on\_exit\ callbacks. A \README.md\ documents that these scripts are for testing purposes only and are not included in the release package.
scripts · high confidence
Behavioural changes
Consumer group rebalancing and heartbeat stability improvements
The consumer group implementation has been refactored to improve stability during rebalances and heartbeat failures. Partition assignment logic is now extracted into a dedicated module, adding retry logic with exponential backoff for topics that are temporarily unavailable (UNKNOWN\_TOPIC\_OR\_PARTITION). Heartbeat handling has been hardened to distinguish between recoverable errors (triggering a rejoin) and terminal errors like authorization failures or fenced instances (triggering a stop). Additionally, a crash-loop bound is now enforced for abnormal heartbeat crashes; if a member crashes too frequently within a sliding time window, it will terminate instead of entering an infinite rejoin loop, helping to prevent cluster-wide thrashing during deployment or outage events.
_lib/kafka\_ex/consumer/consumer\group · high confidence
KafkaEx v1.0 API overhaul with dynamic supervision and telemetry
KafkaEx v1.0 introduces a new primary client API (\KafkaEx.API\) that supports both direct client-pid usage and a \use KafkaEx.API\ mixin for automatic client resolution. The library now uses a \DynamicSupervisor\ to manage client processes, making restart strategies and child limits configurable via \max\_restarts\ and \max\_seconds\. Configuration is centralized in \KafkaEx.Config\, which introduces \:request\_timeout\ (replacing the deprecated \:sync\_timeout\), \:connect\_timeout\, and flexible broker definitions (tuples, CSV strings, or dynamic functions). Additionally, the library now emits comprehensive \:telemetry\ events for requests, connections, authentication, produce/fetch operations, and consumer group lifecycle events to improve observability.
_lib/kafka\ex · high confidence
KafkaEx v1.0 release with new record structures and partitioner behavior
This release introduces the \KafkaEx.Messages.Fetch.Record\ struct to align fetched message metadata with the Java Kafka client's \ConsumerRecord\, including fields like \leader\_epoch\, \timestamp\_type\, and \headers\. It also replaces the previous default partitioning logic with \KafkaEx.Producer.Partitioner.Default\, which uses signed 31-bit murmur2 hashing to match Java client behavior, while the old unsigned 32-bit masking logic is preserved in the deprecated \KafkaEx.Producer.Partitioner.Legacy\ for backward compatibility. Additionally, a new \KafkaEx.API.Behaviour\ is provided to standardize client callback implementations.
_lib/kafka\_ex/api, lib/kafka\_ex/messages/fetch, lib/kafka\ex/producer · high confidence
Major configuration overhaul and migration to Elixir 1.9+ config syntax
The application configuration has been significantly expanded and modernized. The main \config/config.exs\ now uses the \import Config\ syntax (replacing the deprecated \Mix.Config\) and introduces numerous new options for \kafka\_ex\, including \request\_timeout\, \max\_restarts\, \max\_seconds\, \commit\_interval\, \commit\_threshold\, \auto\_offset\_reset\, \sleep\_for\_reconnect\, \metadata\_update\_interval\, and explicit \use\_ssl\/\ssl\_options\. It also adds support for SASL authentication, static consumer-group membership (KIP-345), and configurable broker discovery callbacks. A new \config/test.exs\ file is introduced to configure test-specific behaviors, such as disabling the default worker, setting the \snappy\_module\ for Kayrock, and enabling log capture.
config · high confidence
Migrate Kafka protocol implementations to Kayrock
The KafkaEx protocol layer in \lib/kafka\_ex/protocol\ has been rewritten to use the Kayrock library, introducing explicit version support for ApiVersions (V0–V3), CreateTopics (V0–V5), DeleteTopics (V0–V4), and DescribeGroups (V0–V5). This change ensures correct handling of protocol-specific features such as throttle times, flexible versioning (KIP-482), and new request/response fields, while providing forward compatibility through fallback implementations for future Kayrock versions.
_lib/kafka\ex/protocol · high confidence
New SASL authentication framework with OAuthBearer and MSK IAM support
KafkaEx introduces a comprehensive SASL authentication system supporting four mechanisms: PLAIN, SCRAM (SHA-256/512), OAuthBearer (KIP-255/342), and AWS MSK IAM. The new \KafkaEx.Auth.Config\ module replaces the legacy top-level \:sasl\_username\/\:sasl\_password\/\:sasl\_mechanism\ application environment keys with a unified \:sasl\ map configuration; using the old keys now raises an error to prevent silent unauthenticated connections. The \KafkaEx.Auth.SASL\ orchestrator handles API version negotiation and protocol version selection (including flexible headers v2+), while specific mechanism modules implement the authentication exchange flows.
_lib/kafka\ex/auth · high confidence
New support modules for error handling, retry logic, and dependency validation
This change introduces a new \lib/kafka\_ex/support\ directory containing foundational modules that improve reliability and configuration safety. \KafkaEx.Support.Retry\ centralizes error classification and provides exponential backoff with jitter for transient failures, while \KafkaEx.Support.OptionalDeps\ validates that required optional dependencies (like \snappyer\ or \aws\_signature\) are installed at startup, preventing runtime crashes. Additionally, new exception types in \KafkaEx.Support.Exceptions\ (e.g., \JoinGroupRetriesExhaustedError\) provide detailed context for consumer group failures, and utility modules like \Murmur\ and \VersionHelper\ support internal protocol operations.
_lib/kafka\ex/support · high confidence
Refactor cluster metadata model to support leaderless partitions and coordinator re-discovery
The cluster metadata model has been refactored to introduce dedicated structs for brokers, topics, and partitions, replacing the previous flat representation. This change ensures that partitions without an available leader (leader ID -1) are retained in the metadata rather than being dropped, which prevents partition count mismatches during production. Additionally, the new \ClusterMetadata\ module includes explicit support for tracking consumer group coordinators and provides a \drop\_consumer\_group\_coordinator\ function, enabling the client to correctly re-discover the coordinator when receiving a NOT\_COORDINATOR response.
_lib/kafka\ex/cluster · high confidence
Regenerated SSL certificates and keys
The SSL certificate authority (CA) and server certificates located in the ssl directory have been regenerated. This update replaces the previous local development certificates (issued to 'localhost' by 'KafkaEx') with new key pairs and certificates, ensuring that TLS connections use the current set of credentials.
ssl · high confidence
Renames KafkaEx module and removes legacy consumer implementation
The main application module is now \KafkaEx\ (previously \Kafka\), consolidating application lifecycle and worker management. The legacy \Kafka.Consumer\ module, which handled raw TCP connections and metadata parsing, has been removed in favor of the structured \KafkaEx.API\ and worker-based approaches documented in the new module. This change ensures the public API aligns with the project name and removes deprecated internal implementation details.
lib · high confidence
Fixes
Introduces bounded connection timeouts to prevent indefinite blocking
The network layer now enforces a configurable connect timeout (defaulting to 10 seconds) when establishing TCP or SSL connections to Kafka brokers. Previously, socket creation could block indefinitely if a broker was unreachable, potentially hanging the calling process. This change ensures that connection attempts fail fast with a timeout error rather than waiting for OS-level TCP timeouts, improving application responsiveness during broker outages or network issues.
_lib/kafka\ex/network · high confidence
Test coverage
Added SSL/TLS test fixtures for client and server certificates; Added chaos tests for compression, consumer groups, and network resilience; Added comprehensive test coverage for KafkaEx authentication mechanisms; Added comprehensive test coverage for KafkaEx client retry logic, error handling, and metadata management; Added comprehensive test coverage for KafkaEx core components; Added integration and unit tests for authentication mechanisms and producer partitioning; Added integration tests for JoinGroup protocol versions; Added integration tests for message production reliability, batching, and partitioning; Added network layer tests for socket and client behavior; Added test coverage for Kafka consumer reliability and initialization; Added test coverage for KafkaEx API client operations; Added unit tests for Kafka message structs and protocol helpers; Added unit tests for Kayrock protocol request and response handling; Added unit tests for consumer group stability and static membership; Added unit tests for support modules; Expanded integration test coverage for consumer group stability and fetch behavior; Expanded test infrastructure for consumer group resilience and chaos testing; Expanded unit tests for worker options and new test infrastructure with Mimic and Testcontainers; New integration test suite for KafkaEx lifecycle and API behavior.
Dependencies
139 commits updating dependencies (2 manifests)
A dependency / build maintenance change in (dependencies) — 139 commits (14 fixs), 2 files.
(dependencies) · medium confidence · unverified
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
How this codebase got here
Score
- CAI 67 → 67 (-0.3)
- Rubric changed (rubric-2026.09.15 → rubric-2026.10.1) — scores are not directly comparable.
Lenses
- Code Health 98 → 98 (-0.2)
- Architecture 95 → 85 (-10.1)
- Maturity 64 → 64 (+0.2)
- Readiness 59 → 59 (-0.1)
- Security 78 → 82 (+3.7)
Resolved (6)
- Documentation: no architecture or design documentation (README.md)
- Documentation: no installation or build instructions (README.md)
- Documentation: no usage examples (README.md)
- Further sole-owners (lower concentration)
- Off-boarding risk: anonymized user #1
- Repeated repair: lib/kafka_ex/support/retry.ex (lib/kafka_ex/support/retry.ex)
New (8)
- Duplicated block (5 lines × 2) (lib/kafka_ex/protocol/kayrock/produce/request_helpers.ex)
- Duplicated block (6 lines × 2) (lib/kafka_ex/client/client.ex)
- Duplicated block (6 lines × 4) (lib/kafka_ex/protocol/kayrock/metadata/v4_request_impl.ex)
- Duplicated block (6 lines × 8) (lib/kafka_ex/protocol/kayrock/offset_fetch/any_request_impl.ex)
- Duplicated block (7 lines × 3) (lib/kafka_ex/protocol/kayrock/fetch/v4_request_impl.ex)
- Duplicated block (7 lines × 4) (lib/kafka_ex/protocol/kayrock/produce/v5_response_impl.ex)
- Off-boarding risk: anonymized user #1
- Projects may be oversized for their cohesion
Changes since last survey
- 1 commits — 0 feature/other, 1 fixes
By area
- lib/kafka_ex — 1 commit
Notable commits
- fix: fix(client): extract group-level error atoms so a stale coordinator is refreshed (#599)
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
kafkaex/kafka_ex was measured the same way every project in this corpus was: the same rubric, at a pinned commit, with the result published in full. Point a surveyor at a repository you know and see whether you agree with it.
About this page
- The score is its most recent published measurement, taken on 3 October 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
- Measured at commit d9d6e40445a5a8858d1b05213c007a93d0b5848d — the exact code this score is about.
- Scored under rubric-2026.10.1 — the same rubric and the same method as every other entry in this index.
- Measured by watchdog.canine.dev using codehealth-analyzer preprod-8fe32cd45d00.