Skip to content
CAI
Software that uses CAICheck a score

typelevel/fs2-kafka

58.1

Adequate · 20 September 2026

8.9k

lines of production code

Scala

primary language

1

measurement over time

CAI band scale
CAI lens gauges

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

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.