CleverCloud/pulsar4s
60.4
Adequate · 20 September 2026
3k
lines of production code
Scala
primary language
1
measurement over time
What this system is
This system is a Scala client library for Apache Pulsar that provides asynchronous message production, consumption, and streaming capabilities. It supports multiple effect systems, including Cats Effect, ZIO, Monix, and Scalaz, allowing developers to integrate Pulsar operations into their preferred functional programming contexts. The library also offers automatic serialization for various data formats such as JSON (via Circe, Jackson, json4s, Play JSON, and spray-json) and Avro, alongside integration with Akka and Pekko Streams for reactive data processing pipelines.
How it got here
2018 — Scala 3 migration and effect system integration
21 changes.
The project established a foundational Scala 3 build environment and refactored the core API to be effect-system-agnostic via a generic AsyncHandler. This enabled the addition of native integrations for multiple asynchronous libraries, including Cats Effect, Monix, ZIO, and Scalaz, alongside comprehensive JSON serialization support for Jackson, Circe, Play JSON, spray-json, and json4s. The period also introduced Akka Streams graph stages for Pulsar sources and sinks, accompanied by extensive integration testing across all new modules.
2019–2023 — Integration with ZIO, fs2, and Pekko Streams
7 changes.
This period focused on expanding the library's ecosystem by adding native support for functional programming libraries, including ZIO for async effects, fs2 for functional streams, and Pekko Streams for graph-based processing. It also introduced automatic Avro schema derivation to simplify serialization and added comprehensive test coverage for these new integrations and existing Circe support.
Features
Add Cats Effect integration for asynchronous Pulsar operations
The library now provides a Cats Effect integration via \CatsAsyncHandler\, enabling asynchronous Pulsar client operations (such as producing, consuming, reading, seeking, and transactional processing) to be executed within the \IO\ monad or any other type class instance implementing \cats.effect.Async\. This allows users to compose Pulsar interactions seamlessly with other asynchronous effects in their applications.
pulsar4s-cats-effect/src/main · high confidence
Add JSON4s serialization support for Pulsar schemas
Users can now serialize and deserialize Pulsar messages to and from JSON using the json4s library. This change introduces a new implicit Schema implementation in the pulsar4s-json4s module, allowing any reference type with a Manifest to be automatically converted to JSON bytes for transmission and back from JSON bytes upon receipt.
pulsar4s-json4s/src/main · high confidence
Add Jackson-based JSON serialization support for Pulsar schemas
Users can now serialize and deserialize Pulsar messages to and from JSON using Jackson. This change introduces a new \JacksonSupport\ object that configures an \ObjectMapper\ with Scala module support, relaxed deserialization settings (e.g., ignoring unknown properties, allowing unquoted field names), and specific number serializers. It also provides an implicit \Schema\[T\]\ implementation in the \jackson\ package object, allowing any type with a Manifest to be automatically treated as a BYTES schema that encodes/decodes via the configured Jackson mapper.
pulsar4s-jackson/src/main · high confidence
Add Pekko Streams integration for Pulsar
This release introduces a new \pulsar4s-pekko-streams\ module that provides Pekko Streams graph stages for Pulsar. Users can now consume messages via \source\ (auto-acknowledging), \committableSource\ (manual ack/nack with transaction support), and \sourceReader\ (for non-subscription reading), and produce messages via \sink\ (single topic) and \multiSink\ (multi-topic). The implementation includes the underlying \GraphStage\ implementations and a \package.scala\ exposing the factory methods, along with tests and an example demonstrating the stream wiring.
pulsar4s-pekko-streams · high confidence
Added Play JSON serialization support for Pulsar schemas
The pulsar4s-play-json module now provides an implicit Schema implementation that serializes and deserializes Scala types using Play JSON. Users can now leverage existing Play JSON Reads and Writes instances to handle Pulsar message encoding and decoding without writing custom serialization logic.
pulsar4s-play-json · high confidence
Added Scalaz async handler for Pulsar4s
Users can now integrate Pulsar4s with Scalaz's Task effect type. This change introduces a new ScalazAsyncHandler that bridges Pulsar's asynchronous Java client operations (such as creating producers/consumers, sending messages, seeking, and transaction management) to Scalaz's Task, enabling functional error handling and async composition within Scalaz-based applications.
pulsar4s-scalaz/src/main · high confidence
Added spray-json marshaller for Pulsar schemas
Users can now serialize and deserialize Pulsar messages using spray-json. This new subproject provides an implicit \spraySchema\ that converts any type with available \RootJsonWriter\ and \RootJsonReader\ instances into a Pulsar \Schema\, enabling seamless integration of spray-json for message encoding and decoding.
pulsar4s-spray-json · high confidence
Automatic Avro schema derivation for Pulsar producers and consumers
The pulsar4s-avro module now provides automatic schema derivation for Scala types using avro4s, allowing users to send and receive strongly-typed messages without manually defining Avro schemas. This implementation supports schema versioning by integrating with Pulsar's SchemaInfoProvider to resolve specific schema versions during decoding, and it uses the Avro binary format for efficient serialization.
pulsar4s-avro/src/main · high confidence
Automatic JSON serialization for Pulsar messages via Circe
The pulsar4s-circe module now provides automatic derivation of Pulsar Schema, MessageWriter, and MessageReader instances for any type T that has implicit circe Encoder and Decoder instances. Users can now seamlessly send and receive case classes (e.g., City) as JSON strings in Pulsar messages by simply importing io.circe.generic.auto.\_ and the circe package object, without manually defining schema serialization logic.
pulsar4s-circe/src/main · high confidence
Initial project setup with Scala 3 support and Pulsar 3.3.7 integration
This change establishes the project's foundational configuration, introducing support for Scala 3 via the \.scalafmt.conf\ file and setting up the build environment with \.sbtopts\ for increased JVM memory. It configures publishing to Maven Central via \sonatype.sbt\ and updates the local development environment by setting the Docker Compose Pulsar image to version 3.3.7 with transaction coordination enabled. The README is also expanded to document the client's features, including support for multiple effect types (Future, Monix, Cats Effect, ZIO), Akka Streams integration, and various JSON schema libraries.
(repo-wide) · high confidence
Initial support for functional streams via fs2 integration
The pulsar4s-fs2 module now provides functional stream capabilities using fs2, allowing users to consume and produce Pulsar messages as streams. This includes new \PulsarStreams\ methods for reading from topics (\reader\), consuming from subscriptions in batch or single modes (\batch\, \single\), and writing to topics via pipes (\sink\, \committableSink\). The implementation wraps Pulsar consumers and producers in fs2 \Stream\ and \Pipe\ abstractions, enabling composable, effectful data processing pipelines with automatic resource management for connections.
pulsar4s-fs2 · high confidence
Monix integration for asynchronous Pulsar operations
The Monix module now provides a complete AsyncHandler implementation that bridges Pulsar's Java asynchronous client API with Monix's Task type. This enables users to perform asynchronous operations—including creating producers and consumers, sending messages, seeking, flushing, and managing transactions—using Monix's reactive primitives instead of Scala Futures or blocking calls.
pulsar4s-monix/src/main · high confidence
New Akka Streams graph stages for Pulsar sources and sinks
This change introduces a complete set of Akka Streams graph stages for Pulsar integration, exposing new factory methods in the \pulsar4s.akka.streams\ package. Users can now create auto-acknowledging sources via \source\, manual acknowledgment sources via \committableSource\ (supporting individual ack/nack and transactional contexts), and reader-based sources via \sourceReader\. For publishing, the library provides a standard \sink\ for single-topic production and a new \multiSink\ that accepts \(Topic, ProducerMessage)\ pairs to route messages to multiple topics dynamically. All stages implement a \Control\ interface, allowing users to safely complete, drain, or shut down the underlying Pulsar consumers and producers.
pulsar4s-akka-streams/src/main · high confidence
ZIO integration now supports async operations, transactions, and message reprocessing
The ZIO integration for Pulsar4s has been updated with a new \ZioAsyncHandler\ that enables asynchronous creation of producers, consumers, and readers, as well as async variants for seeking, receiving, and acknowledging messages. Users can now leverage ZIO effects for transaction management (start, commit, abort) and the new \reconsumeLater\ API allows delaying message redelivery. This change shifts the library from synchronous blocking calls to non-blocking ZIO \Task\ effects for all core Pulsar client interactions.
pulsar4s-zio/src/main · high confidence
Behavioural changes
Core API refactored to use an effect-system-agnostic AsyncHandler
The pulsar4s-core library has been rewritten to decouple the client from specific asynchronous implementations. Instead of relying directly on Scala \Future\, the core traits (\Consumer\, \Producer\, \Reader\) and the \PulsarClient\ now operate through a generic \AsyncHandler\[F\[\_\]\]\ type class. This allows users to integrate the library with different effect systems (such as Monix or Scalaz, as indicated by the commit history) by providing their own \AsyncHandler\ instance. The \FutureAsyncHandler\ is provided as the default implementation for \Future\. Additionally, the \MessageId\ wrapper now exposes \batchIndex\, and configuration objects like \ConsumerConfig\ and \ProducerConfig\ have been updated to support newer Pulsar features such as \deadLetterPolicy\, \enableBatchIndexAcknowledgement\, and \batcherBuilder\.
pulsar4s-core/src/main · high confidence
Test coverage
Added Scala 3 integration tests for Circe producer/consumer; Added ZIO integration tests for async operations and transactions; Added initial test structure for json4s module; Added integration tests for consumer metadata, async operations, and transactions; Added tests for Avro schema derivation and producer-consumer round trips; Added tests for Cats-Effect async integration and resource management; Added tests for Jackson schema serialization and deserialization; Added tests for Monix async client operations and transaction support; Added tests for Pulsar Akka Streams components; Added tests for Scalaz async handler integration; Added tests for circe schema derivation and marshalling.
Dependencies
Updated build infrastructure and SBT version
The project's build configuration has been updated to use SBT version 1.12.5, with memory settings configured in .sbtopts to allocate 4GB heap and 1GB metaspace. The plugins.sbt file now includes sbt-pgp 2.3.1, sbt-ci-release 1.11.2, sbt-release 1.4.0, and sbt-sonatype 3.12.2, replacing previous plugin versions and establishing the current release and signing tooling.
project · high confidence
Updated dependency versions in build configuration
The build.sbt file has been updated to use newer versions of key libraries, including Apache Pulsar client 4.0.9, Cats Effect 3.6.3, ZIO 2.1.24, and Scala 3.8.2, alongside updates to Jackson, Log4j, and other transitive dependencies.
(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 60.
Lenses
- Code Health 99
- Architecture 100
- Maturity 52
- Readiness 49
- Security 80
Changes since last survey
- 300 commits — 280 feature/other, 20 fixes
By area
- (root) — 131 commits
- (repo) — 79 commits
- .github/workflows — 26 commits
- pulsar4s-core/src — 22 commits
- project/build.properties — 18 commits
- pulsar4s-akka-streams/src — 6 commits
- project/plugins.sbt — 5 commits
- pulsar4s-cats-effect/src — 5 commits
- pulsar4s-circe/src — 2 commits
- pulsar4s-jackson/src — 2 commits
- pulsar4s-fs2/src — 1 commit
- pulsar4s-pekko-streams/src — 1 commit
- pulsar4s-scalaz/src — 1 commit
- pulsar4s-zio/src — 1 commit
Notable commits
- fix: Bump sbt and fix actions for scala 3
- fix: Fix README.
- fix: Fix deprecated notice in jackson support.
- fix: Fix duplicate github flow after bad rebase.
- fix: Fix nexus url.
- fix: Fix sbt 1.5 syntax.
- fix: Fix test with circe.
- fix: Fix tests
- fix: Fix types for scala 2.x
- fix: Fix workflows and scala 2.12 build
- fix: Fixed async message tests by using strongly typed message attributes (#332)
- fix: Merge pull request #258 from gmethvin/fix-receive-callback
- fix: Merge pull request #364 from valdo404/fix-fs2-publishing
- fix: Override clone due to what seems like a compiler bug
- fix: Override clone due to what seems like a compiler bug
- fix: Upgrade pulsar to 2.7.1 and fix schema issues (#298)
- fix: fix(workflows): fix gpg_private_key input
- fix: fix(workflows): try to get worklows running
- fix: fix(workflows): try to get worklows running, take two
- fix: fix: add fs2 to aggregate
- …and 280 more
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
CleverCloud/pulsar4s 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 8bb8289e04058a619d00d8a3284441d1a5d7d7b8 — 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.