typelevel/fs2-kafka
58.1
Adequate · 20 September 2026
8.9k
lines of production code
Scala
primary language
1
measurement over time
What this system is
This system is fs2-kafka, a functional Scala library that provides stream-based Kafka consumers, producers, and admin clients built on the Cats Effect ecosystem. It manages Kafka interactions through a modular, capability-based API that supports type-safe configuration, secure SSL/TLS authentication, and seamless integration with Avro serialization via Vulcan. The library also includes tooling for automated code migration and comprehensive testing utilities to ensure reliability across different Scala versions.
How it got here
2018 — Typelevel migration and API cleanup
7 changes.
The project migrated to the Typelevel organization, updating Maven coordinates and licensing while overhauling documentation with a new Docusaurus site. Legacy consumer and producer APIs were removed to streamline the codebase, and core dependencies were upgraded to support Scala 3.
2019 — API redesign and robustness improvements
9 changes.
This period focused on a major API overhaul for fs2-kafka 4.x, introducing type-safe settings, new record models, and enhanced offset management. The work also included significant internal refactoring for thread safety, a new configurable commit recovery strategy, and the addition of a Vulcan integration module for Avro serialization.
2020–2026 — Modular API and migration tooling
11 changes.
The project refactored core Kafka components into a modular, trait-based API to improve testability and enable granular capability composition. It introduced automated Scalafix rules to handle breaking changes in type parameters and factory methods, alongside new features for SSL credential management and producer instantiation.
Features
Add Scala 2.12 compatibility layer for Kafka converters
A new compatibility file is introduced in the Scala 2.12 source directory to support Java interop and collection conversions. This file provides a wrapper for \scala.collection.JavaConverters\ and defines an implicit class to convert Scala \Option\ values to Java \Optional\, ensuring the core Kafka functionality works correctly on Scala 2.12.
modules/core/src/main/scala-2.12 · high confidence
Initial Docusaurus-based documentation site
The website has been migrated to Docusaurus, introducing a new documentation structure with a sidebar covering overview, quick example, consumers, producers, transactions, admin, modules, and technical details. The site now features a dedicated footer with license and icon attribution, and includes a link to the API documentation in the header.
website · high confidence
Initial website styling and branding assets
The website now includes custom CSS styles for layout and typography, specifically targeting the navigation footer, main container, and link colors, alongside new SVG logo assets (fs2-kafka) for both light and dark backgrounds.
website/static · high confidence
Introduce MkAdminClient capability for AdminClient instantiation
A new \MkAdminClient\ trait has been added to the \fs2.kafka.admin\ package, providing a capability-based abstraction for creating the underlying Java \AdminClient\. This allows users to override the default instantiation logic (which uses \Sync\ for blocking calls) with custom implementations, such as for testing purposes, by providing an implicit instance in lexical scope.
modules/core/src/main/scala/fs2/kafka/admin · high confidence
Introduce fs2-kafka-vulcan integration module
A new \fs2-kafka-vulcan\ module has been added to provide seamless integration between fs2-kafka and the Vulcan Avro library. This addition introduces \AvroSerializer\ and \AvroDeserializer\ classes that wrap Vulcan codecs, allowing users to serialize and deserialize Avro records using schema registry. The deserializer specifically handles writer schema resolution from the schema registry to ensure compatibility, while both serializers and deserializers are allocated as \cats.effect.Resource\ to manage lifecycle and resources safely.
modules/vulcan/src/main · high confidence
New KafkaCredentialStore for PEM-based SSL configuration
A new \KafkaCredentialStore\ trait and companion object have been added to the \fs2.kafka.security\ package, providing a structured way to manage SSL credentials. The implementation includes a \fromPemStrings\ factory method that accepts CA certificates, client private keys, and client certificates as strings, automatically configuring the necessary Kafka properties (such as \security.protocol\, \ssl.truststore.type\, and \ssl.keystore.type\) to use PEM format. This allows users to easily set up mutual TLS authentication using in-memory PEM data without needing to manage temporary files or external keystores.
modules/core/src/main/scala/fs2/kafka/security · high confidence
Removals
Removal of legacy fs2-kafka consumer and producer APIs
The legacy \KafkaConsumer\, \KafkaProducer\, \ConsumerSettings\, \ProducerSettings\, and related message types (\CommittableMessage\, \ProducerMessage\, \ProducerResult\) have been removed from the \src/main/scala/fs2/kafka\ package. This deletion eliminates the previous actor-based consumer implementation and the \produceWithBatching\ producer API, requiring users to migrate to the updated API surface (such as \KafkaConsumer\#partitionedStream\ and the new \ProducerMessage\/\ProducerResult\ structures) that is being introduced in other parts of this release.
src/main/scala/fs2/kafka · high confidence
Behavioural changes
Add Scala 2.13+ Java interop converters
The library now includes a new internal \converters\ object in the \scala-2.13+\ source directory to handle conversions between Scala and Java collections, options, and durations using the \scala.jdk\ package, which is available in Scala 2.13 and later. This change enables the use of modern Scala collection interoperability features for users running on Scala 2.13+.
modules/core/src/main/scala-2.13+ · high confidence
Automated migration for passthrough parameter reordering in ProducerRecords
A new Scalafix rule named 'Fs2Kafka' has been added to automatically update code affected by the removal of the passthrough type parameter from \ProducerRecords\ and \ProducerResult\. This tool reorders type arguments, moving the passthrough parameter \P\ to the second position (e.g., changing \ProducerRecords\[K, V, P\]\ to \ProducerRecords\[P, K, V\]\), ensuring compatibility with the updated API signatures without requiring manual refactoring.
scalafix/rules · high confidence
Introduce MkProducer capability trait for producer instantiation
The library now exposes a \MkProducer\ trait that abstracts the creation of the underlying Java \KafkaProducer\. This allows users to override the default producer instantiation logic, which is particularly useful for testing scenarios where a mock or stub producer can be provided via an implicit instance instead of the default synchronous implementation.
modules/core/src/main/scala/fs2/kafka/producer · high confidence
New commit recovery strategy and dedicated exception types for offset commits
The library introduces a configurable \CommitRecovery\ strategy for \KafkaConsumer\ that automatically retries offset commits on \RetriableCommitFailedException\ and \RebalanceInProgressException\ using jittered exponential backoff, with a fallback to fixed 10-second intervals before giving up. To support this, new exception types have been added: \CommitRecoveryException\ (raised when retries are exhausted), \CommitTimeoutException\ (a retriable exception for commits exceeding the configured timeout), \ConsumerShutdownException\ (raised for requests after consumer termination), and specific \DeserializationException\ and \SerializationException\ types for serializer/deserializer failures.
modules/core/src/main/scala/fs2/kafka · high confidence
Optional synchronous offset commits on partition revoke
The Kafka consumer now supports committing offsets synchronously when partitions are revoked during a rebalance. This behavior is controlled by a new \commitOnRevoke\ setting; when enabled, the consumer ensures that pending offsets are committed before the rebalance logic proceeds, providing stronger guarantees against data loss in scenarios where immediate consistency is required.
modules/core/src/main/scala/fs2/kafka/internal/actor · high confidence
Project migration to Typelevel organization and documentation overhaul
The project has moved to the Typelevel organization, updating the Maven coordinates from \com.ovoenergy\ to \org.typelevel\ and changing the artifact group. The README has been significantly simplified to serve as a landing page pointing to the new Typelevel microsite for full documentation, while the license file has been updated to remove the Apache License appendix and instead include the MIT license for derived code from the Embedded Kafka library. Additionally, the build configuration has been updated to use Scala Steward for dependency management with pinned versions for Kafka clients and Logback, and the code formatting rules in \.scalafmt.conf\ have been expanded to support Scala 3 source compatibility and stricter style enforcement.
(repo-wide) · high confidence
Redesign of the website home page with updated badges and API docs link
The website home page has been redesigned to improve clarity and provide better access to project resources. The new layout prominently features a link to the API documentation, ensuring users can easily find technical references. Additionally, the status badges have been updated to reflect the project's move to the Typelevel organization, including a switch from Gitter to Discord for community chat, and the version badge now correctly handles version strings containing dashes.
website/pages · high confidence
Refactored KafkaConsumer into a modular trait-based API
The KafkaConsumer implementation has been restructured from a single monolithic class into a set of focused capability traits (KafkaAssignment, KafkaCommit, KafkaConsume, KafkaConsumeChunk, KafkaConsumeGrouped, KafkaConsumerLifecycle, KafkaMetrics, KafkaOffsets, KafkaSubscription, and KafkaTopics). This modular design allows users to compose only the consumer capabilities they need, simplifies testing by enabling easier mocking of specific behaviors, and introduces new features such as chunk-based consumption with automatic offset management (KafkaConsumeChunk), grouped partition streams (KafkaConsumeGrouped), and explicit lifecycle control via terminate and awaitTermination (KafkaConsumerLifecycle).
modules/core/src/main/scala/fs2/kafka/consumer · high confidence
Refactored internal blocking and resource management for Kafka consumers and producers
The internal implementation for Kafka consumer and producer connections has been restructured to improve thread safety and performance. A new \Blocking\ abstraction was introduced to manage synchronous Kafka API calls, allowing users to optionally provide a custom \ExecutionContext\ via \ConsumerSettings\ and \ProducerSettings\ instead of relying on a default single-threaded executor. The \WithConsumer\ and \WithProducer\ components now utilize this abstraction to wrap blocking operations, and \WithConsumer\ exposes a \synchronouslyDuringRebalance\ method to handle offset commits safely during rebalances. Additionally, \WithAdminClient\ was updated to use cancelable futures for better resource management.
modules/core/src/main/scala/fs2/kafka/internal · high confidence
Removal of internal FiniteDuration conversion syntax
The internal \FiniteDurationSyntax\ helper, which provided an \asJava\ method to convert Scala \FiniteDuration\ values to Java \Duration\ objects, has been removed from the codebase. This change eliminates a specific internal utility class used for time unit conversions within the Kafka integration layer.
src/main/scala/fs2/kafka/internal · high confidence
Updated Scalafix rules for factory deprecations and type parameter reordering
The scalafix test suite now includes rules to handle deprecated factory methods and reordering of type parameters. Specifically, the \FactoryDeprecations\ rule migrates code from using \.using()\ on \KafkaProducer\, \KafkaConsumer\, and \TransactionalKafkaProducer\ to the new direct method syntax (e.g., \KafkaProducer\[IO\].resource(settings)\). Additionally, the \PassthroughParams\ rule updates type parameter orderings in \ProducerRecords\, \ProducerResult\, and \TransactionalProducerRecords\ to reflect the new signature where the key type comes before the value type (e.g., \ProducerRecords\[Int, String, String\]\ instead of \ProducerRecords\[String, String, Int\]\).
scalafix/input, scalafix/output · high confidence
fs2-kafka 4.x: Complete API redesign with new type-safe settings and record models
This release introduces a major API overhaul for the core library, replacing the previous configuration and record models with new, type-safe alternatives. Consumer and producer interactions are now configured via dedicated \ConsumerSettings\ and \AdminClientSettings\ classes that expose strongly-typed options (such as \withAutoOffsetReset\ using the new \AutoOffsetReset\ enum, and \withCredentials\ using \KafkaCredentialStore\) instead of raw string properties. Record handling has been redesigned with new \ConsumerRecord\, \ProducerRecord\, and \Headers\ types that include Cats typeclass instances (Eq, Show, Traverse, Bitraverse) for functional composition. Offset management is also restructured, introducing \CommittableOffset\ and \CommittableOffsetBatch\ to replace the previous \CommittableOffset\/\CommittableOffsetBatch\ implementations, enabling more robust batch commits and transactional support via \CommittableProducerRecords\. Additionally, new enums like \Acks\ and \IsolationLevel\ provide explicit control over producer acknowledgments and consumer isolation levels.
repository · high confidence
Test coverage
Added comprehensive test suite for fs2-kafka core components; Added schema compatibility testing utilities for Avro codecs; Added test coverage for fs2-kafka-vulcan integration components; Added test infrastructure for Cats equality and suite base traits; Added test suite for Scalafix rules; Added tests for KafkaCredentialStore.fromPemStrings; Added tests for internal syntax and KafkaFuture cancellation.
Dependencies
Major dependency upgrades and Scala 3 support
This release upgrades core dependencies to their latest versions, including fs2 to 3.14.0, cats-effect to 3.7.1, and kafka-clients to 4.3.1, while adding support for Scala 3 (3.3.8) alongside Scala 2.12 and 2.13. The build configuration also introduces a new scalafix module for testing semantic rewrites and integrates a Docusaurus-based website for documentation.
(dependencies) · high confidence
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
Baseline
- First survey — no prior run to compare against. CAI 58.
Lenses
- Code Health 95
- Architecture 93
- Maturity 45
- Readiness 55
- Security 77
Changes since last survey
- 300 commits — 281 feature/other, 19 fixes
By area
- (root) — 126 commits
- modules/core — 63 commits
- (repo) — 41 commits
- project/build.properties — 20 commits
- project/plugins.sbt — 20 commits
- .github/workflows — 13 commits
- docs/src — 10 commits
- modules/vulcan — 3 commits
- modules/vulcan-testkit-munit — 1 commit
- project/PackagingTypePlugin.scala — 1 commit
- website/pages — 1 commit
- website/siteConfig.js — 1 commit
Notable commits
- fix: Fix CI build issues (#1261)
- fix: Fix MiMa missing class problem
- fix: Fix CallbackStack.Node leak (#1318)
- fix: Fix compilation after #1340 merge (#1353)
- fix: Fix docs publishing (#1252)
- fix: Fix finishing the same Promise multiple times in case of failure (#1393)
- fix: Fix flaky KafkaConsumer#assignmentStream tests (#1274)
- fix: Fix merge issue
- fix: Fix spillover records reset after poll loop (#1398)
- fix: Merge pull request #1197 from fd4s/bplommer/fix-2.x-base-version
- fix: Remove workaround fixed in SBT 1.3
- fix: Revert "Change to publish on series branches"
- fix: Revert "Feat/add parallelization control to consumer" (#1500)
- fix: Revert "Remove deprecated methods (#1268)" (#1296)
- fix: Revert "Update sbt-mdoc to 2.5.1 in series/3.x (#1277)" (#1284)
- fix: Revert #1126 as it's causing a performance regression (#1322)
- fix: Revert accidental slf4j-api upgrade in #1492 (#1503)
- fix: Revert commit(s) d3ffad6f, 2841d5cb
- fix: fix tests
- change: Add 'Reformat with scalafmt 2.2.2' to .git-blame-ignore-revs
- …and 280 more
Architecture
- 0 containers · 1 bounded contexts · 0 dependency edges (baseline)
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
typelevel/fs2-kafka 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 20 September 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
- Measured at commit 59983592fbcb6413e9a7811350003a5a6c09ccf9 — the exact code this score is about.
- Scored under rubric-2026.09.15 — the same rubric and the same method as every other entry in this index.
- Measured by watchdog.canine.dev using codehealth-analyzer preprod-b51f968c9b10.