Skip to content
CAI
Software that uses CAICheck a score

lensesio/stream-reactor

62.9

Adequate · 28 September 2026

79.6k

lines of production code

Scala

with Java

2

measurements over time

CAI band scale
CAI trend line
CAI lens gauges

What this system is

This system is a comprehensive Kafka Connect plugin suite (Stream Reactor) that facilitates data integration between Apache Kafka and a wide variety of external systems, including cloud storage providers, databases, and messaging services. It provides both source and sink connectors for platforms such as AWS S3, Azure, GCP, Cassandra, MongoDB, Elasticsearch, and MQTT, enabling bidirectional data flow. The framework emphasizes reliability through configurable error handling, retry policies, and exactly-once semantics, while supporting complex data transformations via a custom query language and SQL capabilities.

How it got here

2016–2023 — Cloud connector consolidation and new integrations

100 changes.

The project established a shared \kafka-connect-cloud-common\ library to standardize configuration, file handling, and error policies across AWS S3, GCP Storage, and Azure Data Lake connectors. This period also saw the introduction of new connectors for Azure Data Lake, GCP Storage, Elastic 6/7, MongoDB, JMS, and MQTT, alongside a comprehensive migration of the test suite from Java to Scala.

2024 — Java connector expansion and SQL transformation

47 changes.

This period focused on establishing a new Java-based connector infrastructure, introducing the Kafka Connect Query Language (KCQL) parser, and adding SQL-based record transformation capabilities. It also saw the release of several new connectors, including HTTP Sink, Azure Event Hubs, Azure Service Bus, and GCP Pub/Sub, alongside significant enhancements to cloud sink commit policies and configuration standards.

2025–2026 — New connectors and observability expansion

22 changes.

This period focused on introducing new sink connectors for BigQuery, Azure CosmosDB, and OpenSearch, alongside significant enhancements to cloud sink observability and batch processing logic. The work also included substantial expansion of test coverage across existing connectors and the introduction of local development harnesses to support these new integrations.

Features

Add format writer exception and partition hashing utility

The cloud-common module now includes a new FormatWriterException class for handling format-specific write errors and a PartitionHasher utility object that calculates partition indices based on file names and task counts, supporting more predictable data distribution in cloud source connectors.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/formats, kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source/config/distribution · high confidence

Added AvroDataFactory for consistent Avro configuration

A new AvroDataFactory has been introduced in the kafka-connect-common module to centralize the creation of AvroData instances. This factory ensures that enhanced Avro schema support is enabled by default, providing a consistent and properly configured foundation for Avro data handling across the codebase.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/avro · high confidence

Added metrics timing utilities

A new Metrics utility object has been introduced in the kafka-connect-common module, providing two methods for measuring execution duration: a functional version (withTimerF) for effectful code using Cats Effect, and a standard version (withTimer) for regular Scala code. Both methods capture the time taken to execute a block and pass the duration in milliseconds to a provided logging function, enabling users to track performance metrics for specific operations.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/metrics · high confidence

Azure CosmosDB sink connector introduced with bulk mode, proxy support, and configurable throughput

The Azure CosmosDB sink connector is now available, allowing Kafka Connect to stream data into Azure CosmosDB. This implementation supports bulk submission to reduce API chatter, configurable consistency levels, and proxy settings for secure network environments. Users can enable automatic database and collection creation, with the option to provision new collections at a specific manual throughput (defaulting to 400 RU/s). The connector also includes robust error handling policies, retry mechanisms, and configurable flush settings based on size, interval, or count.

kafka-connect-azure-cosmosdb/src/main · high confidence

Azure Data Lake Storage Gen2 sink connector introduced

A new sink connector for Azure Data Lake Storage Gen2 is now available, allowing users to stream Kafka data into ADLS. The connector supports three authentication modes—Default Azure Credential, Shared Key (account name and key), and Connection String—and provides configuration for HTTP timeouts, connection pooling, and error handling policies. It also includes tunable upload performance settings (parallelism, block size, and single-upload size thresholds) and features such as schema change rollover, null record skipping, and configurable local staging areas.

kafka-connect-azure-datalake/src/main · high confidence

Centralized configuration key definitions for cloud connectors

A new \PropsKeyEnum\ has been introduced in the \kafka-connect-cloud-common\ module to serve as the single source of truth for KCQL configuration property keys. This enum defines constants for text reading options (such as regex, start/end tags, and line ranges), envelope storage settings (including key, headers, value, and metadata storage), file naming parameters, flush controls, and post-processing actions (like copy, move, delete, and late arrival handling). By centralizing these keys, the connector ensures consistent configuration handling across source and sink components.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/config/kcqlprops · high confidence

Configurable HTTP retries and timeouts for GCP connectors

GCP connectors now allow users to configure HTTP retry behavior and socket/connection timeouts. The new configuration properties include the number of retries, retry interval, and backoff multiplier, along with optional socket and connection timeout settings in milliseconds. These settings are applied to the underlying GCP service clients to improve reliability and control over network interactions.

java-connectors/kafka-connect-gcp-common/src/main · high confidence

Configurable padding for sink file names

The cloud sink connectors now support configurable padding for the offset and partition fields in output file names. Users can choose between left-padding, right-padding, or no padding via the \padding.strategy\ configuration property, and specify the padding length and character. This ensures consistent file ordering and naming conventions when writing data to cloud storage.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/config/padding · high confidence

Expanded JMX observability for cloud sink operations

The cloud sink now exposes a significantly richer set of JMX metrics via the CloudSinkMetricsMBean, providing granular visibility into writer lifecycle, granular lock caching, garbage collection, and master-lock write dynamics. A new OpTimer class and StorageInterfaceWithMetrics decorator ensure that storage operation latencies and errors are recorded even when the underlying storage SDK throws exceptions, giving operators accurate performance and failure data.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/metrics · high confidence

Initial implementation of GCP Storage Source and Sink connectors

This change introduces the GCP Storage connector, enabling users to ingest data from Google Cloud Storage buckets into Kafka (Source) and write Kafka data to GCS buckets (Sink). The implementation includes the connector and task classes for both directions, along with configuration definitions for error policies, retries, compression, and schema handling. It also provides GCP-specific authentication client creation, bucket name validation, and a directory lister for source polling.

kafka-connect-gcp-storage/src/main · high confidence

Initial implementation of the MQTT Sink connector

This change introduces the core components for the MQTT Sink connector, enabling users to publish Kafka records to an MQTT broker. The new MqttSinkConnector and MqttSinkTask classes handle configuration and task lifecycle, while the MqttWriter class manages the actual message publishing. Key capabilities include dynamic topic targeting (using record keys, topics, or nested struct/map fields), automatic conversion of various data types (Structs, Maps, Collections) to JSON payloads, and support for configurable QoS and retained messages.

kafka-connect-mqtt/src/main/scala/io/lenses/streamreactor/connect/mqtt/sink · high confidence

Initial release of the Kafka Connect Query Language (KCQL) parser and model

This change introduces the core components for the Kafka Connect Query Language, including the ANTLR4 grammar files (ConnectorLexer.g4, ConnectorParser.g4) and the Java model classes (Kcql, Field, FieldType, etc.) that represent parsed queries. It enables users to define data movement and transformation rules using a SQL-like syntax, supporting features such as INSERT/UPSERT modes, field selection and aliasing, partitioning strategies (including NOPARTITION), and target type parsing (key, value, header). The entry also includes example CQL files demonstrating basic usage patterns like ignoring fields and renaming columns.

java-connectors/kafka-connect-query-language/src/main · high confidence

Integrate Graft context graph into Claude and Cursor workflows

This change adds configuration files to integrate the 'Graft' tool into the development environment. For Claude Code, it installs helper scripts and settings that trigger Graft hooks on specific actions (like edits or bash commands) and update the status line, while also granting permissions for Graft-related bash commands. For Cursor, it configures the Graft MCP server and adds a rule file that instructs the AI to use the Graft context graph for understanding code and finding file spans before performing tasks.

.claude, .cursor · high confidence

Introduce Azure Event Hubs source connector

Adds a new Kafka Connect source connector for Microsoft Azure Event Hubs. The connector reads events from Event Hubs and writes them to Kafka topics, supporting configurable KCQL mappings for field selection and data routing. It includes configuration for consumer offset behavior (seeking to earliest or latest), a 30-second close timeout, and validates input/output topic names against a specific pattern. The implementation supports exactly-once processing semantics and handles consumer rebalancing to ensure correct offset management.

java-connectors/kafka-connect-azure-eventhubs · high confidence

Introduce Azure Service Bus Kafka Connect connector

Adds a new Azure Service Bus connector for Kafka Connect, providing both source and sink capabilities. The sink component maps Kafka records to Service Bus messages, supporting configurable retry policies, batch sending, and KCQL-based routing, while the source component ingests Service Bus messages into Kafka topics with configurable prefetch, backoff, and polling behavior.

java-connectors/kafka-connect-azure-servicebus/src/main · high confidence

Introduce BigQuery Sink Connector with batch loading and upsert capabilities

The BigQuery sink connector is now available in the java-connectors/kcbq-connector module, enabling users to stream Kafka data into Google BigQuery. This connector supports standard row-based writes as well as high-throughput batch loading via Google Cloud Storage (GCS), configurable through the \enableBatchLoad\ and \gcsBucketName\ settings. It also introduces support for upserts and deletes, allowing existing records to be updated or removed based on a configured Kafka key field. The connector handles schema evolution automatically, manages GCP authentication via JSON keys or application defaults, and routes failed records to a dead-letter queue for inspection.

java-connectors/kcbq-connector · high confidence

Introduce Elastic 6 sink connector with strict error handling and HTTP basic auth

The Elastic 6 sink connector is now available, providing a dedicated implementation for Elasticsearch 6 clusters. This connector supports HTTP basic authentication and SSL/TLS via key/trust stores, and enforces document types in index, update, and delete operations (defaulting to the index name when not explicitly configured). Per-item bulk errors are surfaced by default (strict mode), allowing the connector to distinguish between 429 rate-limit errors and mapper errors; a tolerant mode is available to log and drop these errors instead. The connector also uses a 5-minute write timeout for bulk operations and index creation.

kafka-connect-elastic6/src/main/scala/io/lenses/streamreactor/connect/elastic6 · high confidence

Introduce FTP/SFTP source connector with file slicing and offset tracking

The FTP source connector now supports monitoring FTP, FTPS, and SFTP servers to stream file contents to Kafka topics. It introduces a configurable file-slicing mechanism (via the \connect.ftp.monitor.slicesize\ setting) that allows large files to be processed in chunks rather than requiring a full download, and tracks file offsets to resume incomplete reads or detect updates. The connector provides two file conversion modes—\MaxLinesFileConverter\ for line-delimited files and \SimpleFileConverter\ for whole-file records—and supports both string and struct key styles for the resulting Kafka source records.

kafka-connect-ftp/src/main/scala · high confidence

Introduce GCP Pub/Sub Source Connector

Adds a new GCP Pub/Sub Source Connector that ingests messages from Google Cloud Pub/Sub subscriptions into Kafka topics. The connector supports configurable output modes (DEFAULT and COMPATIBILITY) to control how message keys, values, and headers are mapped to Kafka Connect records, and includes validation to prevent configuration conflicts when multiple KCQL statements reference the same subscription.

java-connectors/kafka-connect-gcp-pubsub/src/main · high confidence

Introduce HTTP Sink connector for Kafka Connect

Adds a new HTTP Sink connector that sends Kafka records to an HTTP endpoint. The connector supports configurable batching (by count, size, or time interval), multiple HTTP methods (POST, PUT, PATCH), and authentication modes including Basic and OAuth2. It includes built-in retry logic with fixed or exponential backoff, backpressure via a semaphore-based queue, and metrics reporting for both successes and failures.

kafka-connect-http/src/main · high confidence

Introduce OpenSearch Sink Connector with secure authentication and strict error handling

Adds a new Kafka Connect sink connector for OpenSearch 2.x, providing a dedicated implementation (KOpenSearchClient) that validates cluster compatibility at startup to prevent misconfiguration against Elasticsearch or unsupported 1.x clusters. The connector supports AWS SigV4 signing and JWT Bearer token authentication (with file-backed rotation and path-traversal guards), while enforcing a cleartext-auth guard by default to prevent credential interception. It also introduces a strict bulk item error mode (enabled by default) that surfaces per-item failures to the error policy, replacing the previous tolerant behavior that silently dropped failed documents.

kafka-connect-opensearch/src/main · high confidence

Introduce SQL-based record transformation in Kafka Connect

Added a new SQL transformation capability to the Kafka Connect SQL common module, allowing users to filter, project, and flatten data in message keys and values using SQL queries. This change introduces the \Transformation\ class, which reads SQL configuration per topic, and a suite of supporting components including \Sql\ for parsing MySQL-compatible queries via Apache Calcite, \StructSql\ for applying transformations to Kafka Connect Structs, \Transform\ for handling various payload types (JSON, Struct, Bytes), and \FieldValueGetter\ for navigating nested schema paths. Users can now configure SQL statements to reshape incoming and outgoing records directly within the connector configuration.

kafka-connect-sql-common/src/main/scala/io/lenses/connect · high confidence

Introduce configurable Kafka-based sink reporting

The sink-reporting module now supports sending connector success and error reports to a configurable Kafka topic. This change introduces a new reporting pipeline (ReportHolder, ReportSender, ReportingController) that uses a dedicated Kafka producer, allowing users to enable reporting and specify the target topic, partition, and security settings (SASL/SSL) via connector configuration prefixes (connect.reporting.error.config. and connect.reporting.success.config.). The implementation ensures reliable lifecycle management by properly shutting down the reporting executor and producer on connector close to prevent thread leaks.

java-connectors/kafka-connect-sink-reporting/src/main/java/io/lenses/streamreactor/connect/reporting · high confidence

Introduce configurable file naming with template support and validation

The cloud sink connectors now support a new \TemplateFileNamer\ that allows users to define custom object key patterns via the \object.key.template\ KCQL property. This feature introduces placeholders for topic, partition, offsets, timestamps, record count, and extension, enabling precise control over file organization in cloud storage. To prevent data loss from silent overwrites, the system enforces validation at configuration time, rejecting templates that lack at least one flush-varying placeholder (such as \{start-offset}\ or \{record-count}\) or that contain unknown placeholders. This logic is implemented in the \kafka-connect-cloud-common\ module alongside refactored naming traits and classes like \CloudKeyNamer\ and \FileNamer\.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/naming · high confidence

Introduce shared elastic-common module for Elastic and OpenSearch connectors

A new \kafka-connect-elastic-common\ library has been added to centralize shared logic for the Elastic and OpenSearch sink connectors. This module provides core utilities for transforming Kafka Connect data into JSON (\JsonPayloadExtractor\, \Transform\), extracting primary keys from values, keys, and headers (\PrimaryKeyExtractor\, \TransformAndExtractPK\), and managing bulk operations (\KBulkClient\, \BulkItemErrorClassifier\). It also standardizes configuration handling (\ElasticCommonConfigDef\, \ElasticCommonSettings\) and defines common constants, enabling consistent behavior and easier maintenance across both connector implementations.

kafka-connect-elastic-common · high confidence

Introduces configurable commit mode for PARTITIONBY sinks

Adds a new \CommitMode\ configuration option for PARTITIONBY cloud sinks, allowing users to switch between the existing \Granular\ mode (default, where each writer commits independently against its own granular lock) and a new \Batch\ mode (where all open writers on a topic-partition commit together through a single master-lock compare-and-swap). This change introduces the \CommitMode\ enum, the \IndexManager\ trait defining the interface for both modes, and the \IndexManagerV2\ implementation which handles the underlying locking, caching, and deduplication logic (including legacy lock migration from granular to batch). Users can now configure their sink to use batch commits to improve exactly-once semantics when partition keys are not deterministic, though switching modes requires stopping the connector first.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/seek · high confidence

Introduces configurable commit policies for cloud sink flushing

Adds a new commit policy framework in the cloud sink common module that determines when a partition file should be flushed and made visible. This framework supports configurable conditions based on record count, file size, and time intervals since the last flush, allowing users to control when data is committed to cloud storage.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/commit · high confidence

Introduces new SQL-focused configuration and conversion utilities

This change adds a new set of Scala classes and traits to the \kafka-connect-sql-common\ module to support SQL-based data projection and conversion. Specifically, it introduces \KcqlWithFieldsSettings\ and \Projections\ to manage KCQL settings such as field mappings, primary keys, and write modes. It also adds several converter utilities, including \MsgKey\ for handling message keys, \AvroConverter\ and \BytesConverter\ for sink record transformations, and helper classes like \SinkRecordConverterHelper\ and \StructHelper\ for extracting and reducing schema fields. These components provide the foundational logic for processing and converting data within the SQL connector context.

kafka-connect-sql-common/src/main/scala/io/lenses/streamreactor/common · high confidence

Introduces schema-driven KCQL property configuration

Adds a new type-safe configuration layer for KCQL properties in the common module, allowing connectors to define and validate settings using explicit schemas. This change introduces \KcqlProperties\ and \KcqlPropsSchema\ classes that enforce data types (such as Int, Long, Boolean, Char, and Map) and support case-insensitive key matching, providing a more robust and structured way to handle connector configuration options compared to previous ad-hoc parsing.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/config · high confidence

Introduces staged batch commit mode for cloud sinks

The cloud sink connectors now support a batch commit mode that stages writer output to temporary paths before atomically moving them to their final destination via a master-lock compare-and-swap. This change introduces new state management classes (CommitState, WriteState) and a dedicated commit manager (WriterCommitManager) to handle the staging, batching, and atomic finalization of files, improving durability and reducing latency by avoiding intermediate copy-and-delete operations.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/writer · high confidence

Introduction of domain model for Kafka topic partition offsets

A new set of case classes (Topic, Offset, TopicPartition, TopicPartitionOffset) has been added to the common model package to represent Kafka topic partition offsets. These classes provide a structured, type-safe way to handle topic names, partition numbers, and offsets, including implicit ordering and conversion methods to and from Kafka's native TopicPartition class. This change introduces new data structures for tracking offsets, which will be used by components that need to manage or backup multiple topics at once.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/model · high confidence

MongoDB sink connector implementation

The MongoDB sink connector is now available, enabling Kafka Connect to write records to MongoDB. This change introduces the core connector components, including MongoSinkConnector and MongoSinkTask for lifecycle management, MongoWriter for handling bulk insert and upsert operations, and utilities like KeysExtractor and SinkRecordToDocument for converting Kafka Connect schema and data structures into MongoDB-compatible documents.

kafka-connect-mongodb/src/main/scala/io/lenses/streamreactor/connect/mongodb/sink · high confidence

New CloudLocation model and staging file utilities in cloud-common

This change introduces the CloudLocation model and associated utilities in the cloud-common module. CloudLocation now represents a specific location in cloud storage with support for bucket, prefix, path, line number, and timestamp, including validation via CloudLocationValidator and helper methods for navigation. Additionally, FileUtils provides staging file operations with a configurable write buffer size (defaulting to 64 KB), enabling more efficient local staging during cloud data processing.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/model/location · high confidence

New MQTT Source Connector implementation

This change introduces the core Scala implementation for the MQTT source connector, adding the MqttSourceConnector, MqttSourceTask, MqttManager, and MqttSSLSocketFactory classes. The MqttSourceTask handles configuration validation (including SSL certificate file existence checks) and instantiates the MqttManager to manage the MQTT client connection. The MqttManager is responsible for subscribing to MQTT topics, matching incoming messages to configured KCQL statements (supporting wildcards and regex), converting payloads via registered converters, and queuing SourceRecords for Kafka Connect. The MqttSSLSocketFactory provides SSL/TLS support using BouncyCastle for loading CA, client certificates, and private keys. The MqttSourceConnector manages task configuration splitting, including logic to replicate shared subscription topics across all tasks if configured.

kafka-connect-mqtt/src/main/scala/io/lenses/streamreactor/connect/mqtt/source · high confidence

New SQL-based JSON transformation capabilities in Kafka Connect SQL Common

The \kafka-connect-sql-common\ module now includes a new \io.lenses.json.sql\ package containing \JacksonJson\ and \JsonSql\ utilities. \JacksonJson\ provides a configured Jackson \JsonMapper\ that uses \NON\_EMPTY\ inclusion to exclude empty values during serialization. \JsonSql\ introduces extension methods on Jackson \JsonNode\ objects, allowing users to execute SQL \SELECT\ queries (parsed via Apache Calcite) against JSON data. This includes support for flattening JSON structures into a single-level object and handling specific field selections, enabling SQL-style transformations directly on JSON payloads within the connector framework.

kafka-connect-sql-common/src/main/scala/io/lenses/json · high confidence

New SimpleJsonConverter for Kafka Connect JSON serialization

A new JSON converter implementation has been added to the Kafka Connect common library, enabling the conversion of Kafka Connect data structures into JSON format. This feature introduces a \SimpleJsonConverter\ class and supporting \SimpleJsonNode\ utilities that handle the serialization of various data types, including primitives, dates, timestamps, decimals, arrays, maps, and structs, into JSON nodes using Jackson. Users can now utilize this component to serialize Connect data into JSON for storage or transmission.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/json · high confidence

New cloud source task infrastructure with late-arrival handling and partition discovery

The \kafka-connect-cloud-common\ module now includes the core \CloudSourceTask\ implementation, introducing a background \LateArrivalTouchTask\ that periodically updates file timestamps to ensure late-arriving data is processed, alongside a new \CloudPartitionSearcher\ for discovering new partitions in cloud storage. This change also adds support for writing watermark information (partition and offset) to Kafka headers via a new configuration flag, and introduces a unified \SourceWatermark\ system for managing source offsets and context.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source · high confidence

New common cloud sink format writers and schema change detection

The \kafka-connect-cloud-common\ module now includes a new set of format writers (Avro, Parquet, JSON, CSV, Text, and Bytes) that handle writing Kafka records to cloud storage. These writers support configurable compression codecs (such as Snappy, GZIP, ZSTD, and Brotli) and introduce schema change detection strategies—allowing users to choose between version-based or compatibility-based detection to determine when to roll over output files.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/formats/writer · high confidence

New common configuration, security, and utility infrastructure for connectors

This change introduces a foundational set of classes in the Kafka Connect common library to standardize how connectors handle configuration, security, and operational utilities. It adds a new \ConfigSource\ abstraction and \BaseConfig\ base class to decouple configuration parsing from specific connector implementations, alongside \KcqlSettings\ for unified KCQL parsing and \RetryConfig\ for standardized retry behavior. Security is enhanced with new \StoresInfo\, \KeyStoreInfo\, and \TrustStoreInfo\ classes that manage SSL/TLS context initialization and keystore loading. Additionally, utility classes like \TopicPartitionOffsetAndMetadataStorage\ for offset management and \EitherUtils\ for functional error handling are added to support these new patterns.

java-connectors/kafka-connect-common/src/main · high confidence

New common utility libraries for Kafka Connect Cassandra

This change introduces a set of new shared Scala utilities within the \kafka-connect-cassandra\ connector to standardize common operations. Users benefit from improved concurrency handling via \ExecutorExtension\ and \FutureAwaitWithFailFastFn\, which provide safer ways to execute tasks and wait for futures with fail-fast semantics. Configuration is streamlined with \ThreadPoolSettings\ for automatic thread pool sizing and \SSLConfigContext\ (now deprecated in favor of \StoresInfo\) for managing SSL contexts. Additionally, \OffsetHandler\ simplifies offset recovery for source tasks, and \QueueHelpers\ offers utilities for draining queues in batches.

kafka-connect-cassandra/src/main/scala/io/lenses/streamreactor/common · high confidence

New configuration schema for cloud source connectors

The cloud source connectors now support a structured configuration schema for post-processing actions (Delete or Move), text reading modes (Regex, Start/End Tag, or Start/End Line), and late-arrival watermarking. This change introduces the \CloudSourcePropsSchema\ and associated enumerations in the \kafka-connect-cloud-common\ module, enabling users to define these behaviors via Kcql properties for more granular control over file processing and data ingestion.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source/config/kcqlprops · high confidence

New configuration validation and Avro serialization utilities

The \kafka-connect-common\ module now includes \Helpers.scala\ to validate that Kafka Connect input topics match the topics defined in KCQL configurations, throwing a \ConfigException\ if there is a mismatch. Additionally, \AvroSerializer.scala\ provides utilities for serializing and deserializing Avro records, including methods to write to streams, read from streams, and convert between generic records and typed products.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/common/config · high confidence

New converters for byte, text, and envelope data formats

Added five new converter implementations in the cloud-common module to handle diverse input formats for cloud source connectors. BytesOutputRowConverter and TextConverter now support raw byte and string payloads respectively, while SchemaAndValueConverter handles structured schema-and-value pairs. Additionally, SchemaAndValueEnvelopeConverter and SchemalessEnvelopeConverter enable ingestion of complex envelope structures (containing key, value, headers, and metadata) for both schema-aware and schemaless (JSON) data, with support for extracting metadata fields like timestamp and partition from the payload.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/formats/reader/converters · high confidence

New local repro harness for GCS sink time-bucket partitioning and local smoke test for OpenSearch sink

Added a new \dev-scripts/gcs-local\ harness that spins up a local Kafka and Connect stack to reproduce and demonstrate how fine-grained \HHmm\ partition keys cause directory proliferation in the GCS sink, and how applying a 15-minute rolling window via Lenses SMTs collapses these into manageable buckets. The harness includes three connector configurations (\gcs-bad\, \gcs-good\, \gcs-rolling\) and scripts to build, deploy, produce data, and verify the resulting GCS directory structure. Additionally, added a local smoke test environment for the OpenSearch sink connector (\dev-scripts/opensearch-local\) that sets up a local Kafka, Connect, and OpenSearch stack with TLS and basic auth to exercise the connector end-to-end.

dev-scripts · high confidence

New schema-aware and schemaless envelope transformers for cloud sinks

The cloud sink connector now includes a new set of message transformers in the common module to handle envelope creation for different data formats. For JSON sinks, a \SchemalessEnvelopeTransformer\ wraps messages into a map-based structure, while for Avro and Parquet sinks, an \AddConnectSchemaTransformer\ and \EnvelopeWithSchemaTransformer\ pair ensures that payloads without native Connect schemas (such as arrays of maps) are enriched with proper schema definitions before being wrapped in an envelope containing key, value, headers, and metadata. This change ensures consistent envelope structures across all supported sink formats.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/transformers · high confidence

New sink record extraction and conversion utilities

This change introduces a new set of components in the \kafka-connect-common\ module to support extracting and converting data from Kafka Connect sink records. It adds a \SinkData\ type hierarchy and a \ValueToSinkDataConverter\ to normalize record values and headers into a unified internal representation. Additionally, it provides a suite of extractors (\KafkaConnectExtractor\, \StructExtractor\, \MapExtractor\, \ArrayExtractor\, etc.) that allow users to navigate and extract specific fields from complex nested structures (structs, maps, arrays) using dot-notation paths, including support for quoted segments to handle field names containing dots.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/extractors · high confidence

New source record converters for Avro and JSON payloads

This change introduces a new set of source record converters in the \kafka-connect-sql-common\ module, providing dedicated implementations for Avro (\AvroConverter\) and various JSON formats (\JsonSimpleConverter\, \JsonOptNullConverter\, \JsonPassThroughConverter\, \JsonResilientConverter\, and \JsonConverterWithSchemaEvolution\). These converters implement a common \Converter\ trait to transform raw byte payloads into Kafka Connect \SourceRecord\ objects, handling schema resolution, key extraction, and data type mapping. The addition of \JsonResilientConverter\ specifically improves robustness by ignoring malformed JSON messages rather than failing, while \JsonConverterWithSchemaEvolution\ supports dynamic schema updates. A \KeyExtractor\ utility is also included to support extracting specific fields from JSON nodes and Connect structs for use as record keys.

kafka-connect-sql-common/src/main/scala/io/lenses/streamreactor/connect · high confidence

New state management and reader builder for cloud source connectors

The \kafka-connect-cloud-common\ module introduces \CloudSourceTaskState\ and \ReaderManagerBuilder\ to centralize the lifecycle and construction of source readers. \CloudSourceTaskState\ manages a collection of \ReaderManager\ instances, handling polling for records, committing watermarks, and graceful shutdown. \ReaderManagerBuilder\ constructs these readers, supporting features like exponential backoff for empty sources, context-aware resumption from previous states, and configurable watermark header injection. This refactoring provides a reusable foundation for cloud source connectors to manage file discovery and data ingestion more robustly.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source/state · high confidence

New test infrastructure for OpenSearch security and various service containers

The test-common module now includes comprehensive support for testing OpenSearch with security enabled, featuring PKI certificate generation, TLS configuration, and security plugin setup (basic auth, client-cert, and JWT). Additionally, new test container wrappers have been added for Azurite, Cassandra, Cosmos DB, GCP Storage, MongoDB, and Redis, along with model classes for Order and blockchain transactions to support integration testing across these services.

test-common · high confidence

New text input readers for structured record parsing

The \kafka-connect-common\ module now includes a suite of new text readers that allow connectors to parse input streams based on specific structural patterns rather than simple line-by-line processing. Users can now utilize \LineStartLineEndReader\ to extract records delimited by start and end markers, \PrefixSuffixReader\ to capture data between custom prefix and suffix strings, \RegexMatchLineReader\ to filter lines matching a regular expression, and \SequenceBasedLineReader\ to identify records using sequential line numbering. These components, along with supporting utilities like \LineReader\ and \OptionIteratorAdaptor\, provide the foundation for handling complex, multi-line text formats in source connectors.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/io · high confidence

New utility classes added to cloud-common library

The \kafka-connect-cloud-common\ module now includes several new utility components to support connector operations: \BytesOutputRow\ provides a wrapper for byte array record values; \IteratorOps\ adds safe skipping logic for iterators; \MapUtils\ offers a method to merge connector context and start properties; \PollLoop\ introduces functional polling and retry mechanisms using Cats Effect; and \TimestampUtils\ simplifies parsing epoch millisecond timestamps into \Instant\ objects.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/formats/bytes, kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/utils · high confidence

New utility classes for type conversions, error handling, and progress tracking

Added several new utility components to the common utils package: CyclopsToScalaEither and CyclopsToScalaOption provide conversion methods between Cyclops and Scala standard library types; EitherOps introduces implicit extension methods to unpack or throw exceptions from Either results; JarManifestProvided offers a trait to easily access JAR manifest version information; and ProgressCounter enables logging of record delivery progress per topic at configurable intervals.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/common/utils · high confidence

Project initialization and repository configuration

The repository has been initialized with essential configuration files and documentation. This includes a \.gitignore\ to manage build artifacts and IDE files, a \.jvmopts\ file to set JVM memory and stack size limits, and build tool configurations for Scala formatting (\.scalafmt.conf\) and import organization (\.scalafix.conf\). The project also introduces a \.mcp.json\ file to configure the 'graft' MCP server, a \suppression.xml\ to handle false-positive security warnings from the OWASP Dependency Check tool, and a \WORKFLOW.md\ guide explaining the GitHub Actions module-based CI/CD pipeline. Additionally, the \README.md\ has been updated to reflect the Lenses Connectors branding, support policies, and build instructions.

(repo-wide) · high confidence

Streaming Parquet input support for cloud storage

The \kafka-connect-cloud-common\ module now includes a new streaming Parquet reader implementation (\ParquetSeekableInputStream\ and \ParquetStreamingInputFile\) that allows Parquet files to be read from cloud storage without loading the entire file into memory. This change enables more efficient processing of large Parquet datasets by supporting chunked reads and stream recreation, addressing previous limitations in handling large files via S3.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/formats/reader/parquet · high confidence

Architecture

Establishes Java connectors build infrastructure and code style standards

The java-connectors module is initialized with a Gradle build system (version 8.14) and comprehensive code quality configurations. This includes standardized code formatting rules for IntelliJ IDEA and Eclipse, Checkstyle linting rules based on Google Java Style, and license header templates for source files. These changes provide the necessary tooling and style guidelines for developing and maintaining Java-based connectors within the project.

java-connectors · high confidence

Behavioural changes

AWS S3 Connector configuration and client initialization refactored

The AWS S3 connector's configuration and client initialization have been restructured to improve clarity and maintainability. A new \AwsS3ClientCreator\ now handles the construction of the AWS SDK S3 client, centralizing logic for authentication (supporting both explicit credentials and default providers), region selection, custom endpoints, and HTTP connection pooling. Configuration properties have been standardized under the \connect.s3\ prefix, with a new \DeprecationConfigDefProcessor\ that enforces the migration from legacy property names (e.g., \aws.access.key\ to \connect.s3.aws.access.key\) by failing startup if old keys are detected. Additionally, a new \DeleteModeSettings\ configuration allows users to choose between \BatchDelete\ and \SeparateDelete\ modes for cleaning up index files, and the connector now includes dedicated ASCII art headers for both the sink and source connectors to aid in log identification.

kafka-connect-aws-s3/src/main · high confidence

Cassandra connector implementation refactored to use JSON-based writes and new connection handling

The Cassandra connector's data path has been restructured: the sink now uses a new \CassandraJsonWriter\ that writes records to Cassandra using the \INSERT ... JSON\ syntax, replacing the previous approach, and connection setup is centralized in a new \CassandraConnection\ object that handles load balancing policies, timeouts, and SSL. Configuration is now defined in \CassandraConfig\ and \CassandraConfigConstants\, and settings are parsed into \CassandraSourceSetting\ and \CassandraSinkSetting\ via \CassandraSettings\, enabling features like configurable fetch size, consistency levels, and row deletion in the sink.

kafka-connect-cassandra/src/main/scala/io/lenses/streamreactor/connect · high confidence

Cloud sink configuration refactored with new validation and retry settings

The cloud sink connector's configuration handling has been restructured into a modular set of settings traits (e.g., FlushSettings, CommitRetrySettings, SkipNullSettings) that are mixed into the main config builder. This change introduces configurable commit-retry parameters (max attempts, base delay, multiplier, max delay) to handle transient network errors during file copy and delete operations, and adds a validator that rejects non-deterministic wallclock-based PARTITIONBY keys when exactly-once semantics are enabled in granular commit mode to prevent silent data loss. Additionally, the legacy KCQL keywords WITH\_FLUSH\_COUNT/SIZE/INTERVAL and WITHPARTITIONER are now rejected in favor of PROPERTIES-based configuration, and a new validation ensures that key suffixes do not start with digits to avoid filename collisions.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/config · high confidence

Cloud sink connectors now support configurable error policies (RETRY, NOOP, THROW) with integrity-aware handling

The cloud sink connectors now allow operators to configure how transient and fatal errors are handled via a new error policy setting. When an error occurs, the connector classifies it as fatal, non-fatal, or integrity-sensitive. Fatal errors always fail the task immediately. Non-fatal errors can be swallowed (logged and ignored) or retried depending on the policy. Crucially, integrity-sensitive errors (such as those affecting exactly-once semantics or lock state) are never silently swallowed by the NOOP policy; instead, they trigger a \RetriableIntegrityException\ that forces a task failure or retry, preventing data loss. This change introduces new exception types (\FatalCloudSinkError\, \NonFatalCloudSinkError\, \BatchCloudSinkError\) and a \RetriableIntegrityException\ to enforce this behavior, ensuring that critical errors are always surfaced while allowing more lenient handling for truly transient issues.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink · high confidence

Cloud storage connectors gain unified partition search, file filtering, and resilient commit-chain retries

Cloud source connectors now use a unified partition search mechanism that walks directory trees to a configured depth, ensuring consistent behavior across S3, GCS, and ADLS, and support filtering files by extension. Cloud sink connectors are more resilient to transient network issues, as the commit-chain Copy and Delete operations now use bounded exponential-backoff retries for errors like TCP resets, while correctly aborting retries on thread interruption to allow prompt shutdown. Additionally, the storage layer introduces a new exact-key check for file existence, a mechanism to 'touch' late-arrival files to update their timestamps, and improved error handling for empty or missing lock files to aid in crash recovery.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/storage · high confidence

Elastic 6 connector configuration model restructured

The configuration handling for the Elastic 6 sink connector has been reorganized into a new set of files (ElasticConfig, ElasticConfigConstants, and ElasticSettings). This change introduces explicit configuration definitions for connection details (protocol, hosts, port, cluster name, prefix), error handling policies (noop, throw, retry with configurable retries and intervals), and performance settings (batch size, write timeout). It also adds support for HTTP Basic Authentication credentials, a configurable primary key joiner separator, and a strict mode for bulk item errors, ensuring that per-item bulk errors are surfaced rather than dropped.

kafka-connect-elastic6/src/main/scala/io/lenses/streamreactor/connect/elastic6/config · high confidence

Elastic 7 connector configuration and settings initialization

The Elastic 7 sink connector now initializes its configuration via a dedicated \ElasticConfig\ definition and \ElasticSettings\ factory. This change introduces explicit support for HTTP Basic Authentication (username and password), configurable write timeouts, batch sizes, and retry intervals, while also enabling strict bulk item error reporting to surface per-item failures rather than dropping them silently.

kafka-connect-elastic7/src/main/scala/io/lenses/streamreactor/connect/elastic7/config · high confidence

HTTP sink batching now correctly packs records using time-based intervals

The HTTP sink's batching logic has been rewritten to fix a defect where time-only batching (configured via \connect.http.time.interval\) caused the writer to send one HTTP request per record instead of packing them into batches. The new implementation introduces a \BatchPolicy\ with composable conditions (count, file size, and interval) that are evaluated using an AND logic for capacity limits and OR logic for flush triggers, ensuring that an elapsed interval acts as a greedy trigger rather than a rejection condition. This change, along with the introduction of \BatchInfo\ and \HttpCommitContext\ structures in the common batch module, ensures that records are correctly grouped and flushed, significantly improving throughput for time-based batching scenarios.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/common/batch · high confidence

Introduce Elastic 7 sink connector with strict bulk error handling

This change introduces the Elastic 7 sink connector components (connector, task, writer, and HTTP client) which now surface per-item bulk errors by default. Previously, individual item failures in a bulk write were silently dropped in tolerant mode; the new implementation sets strict item errors to true, causing the connector to report errors (including classification of 429 vs mapper errors) so users are aware of partial failures rather than having records silently lost.

kafka-connect-elastic7/src/main/scala/io/lenses/streamreactor/connect/elastic7 · high confidence

Introduce common cloud data format readers

The \kafka-connect-cloud-common\ module now includes a suite of new stream readers for Avro, JSON, CSV, Parquet, and raw bytes, providing a shared foundation for cloud source connectors. This addition includes specific behavioral fixes: Avro and Parquet readers now disable Avro's fast-read mode to prevent DATE logical-type fields from being incorrectly cast to \java.time.LocalDate\ instead of the expected raw integer, and the JSON reader adds support for GZIP-compressed streams.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/formats/reader · high confidence

Introduces configurable commit policies and partition path validation for cloud sinks

The cloud sink implementation now includes a default commit policy that triggers flushes based on file size, time interval, and record count, alongside a new configuration utility for validating and managing partition name paths. This ensures that sink operations adhere to defined thresholds for data persistence and that partition identifiers are correctly formatted without reserved characters.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/commit, kafka-connect-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/config · high confidence

Introduces new model classes for structured Kafka reporting records

The reporting module now includes a new set of model classes to structure data sent to the reporting topic. A new \ReportingRecord\ class encapsulates standard metadata (topic, partition, offset, timestamp, endpoint, payload) alongside connector-specific data via a generic \ConnectorSpecificRecordData\ interface. A \RecordConverter\ handles the transformation of these records into Kafka \ProducerRecord\s, explicitly mapping fields to headers defined in \ReportHeadersConstants\ (such as \input\_offset\, \input\_timestamp\, etc.), replacing previous ad-hoc header handling.

java-connectors/kafka-connect-sink-reporting/src/main/java/io/lenses/streamreactor/connect/reporting/model · high confidence

Introduces structured schema for cloud sink KCQL properties

The cloud sink configuration now uses a formalized schema (SinkPropsSchema) to define and validate KCQL properties. This change explicitly maps configuration keys such as ObjectKeyTemplate, FlushCount, and various StoreEnvelope flags to their specific data types, ensuring that properties like padding, partitioning, and envelope storage are parsed and validated consistently when reading from KCQL statements.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/config/kcqlprops · high confidence

JMS connector refactored to use Jakarta JMS API and new configuration system

The JMS connector implementation has been rewritten to migrate from the legacy Java EE JMS API to the Jakarta JMS API (imported as \jakarta.jms\), ensuring compatibility with modern JMS providers. This change introduces a new configuration and settings layer (\JMSConfig\, \JMSSettings\, \JMSConfigConstants\) that centralizes connector properties and introduces a flexible converter system (\ConverterConfigurator\, \ConverterClassLoader\) for handling source and sink message transformations. The core session management is now handled by a new \JMSSessionProvider\ which manages JMS connections, sessions, and destination lookups (supporting both JNDI and CDI selectors) for both queues and topics.

kafka-connect-jms/src/main/scala · high confidence

MQTT connector configuration and connection logic refactored

The MQTT connector's configuration and connection handling have been restructured to use a new, centralized settings model. The \MqttConfig\ object now explicitly defines all connector parameters (such as hosts, credentials, QoS, SSL paths, and KCQL) using Kafka Connect's \ConfigDef\, while dedicated traits and case classes (\HostSettings\, \MqttSourceSettings\, \MqttSinkSettings\) manage the parsing and validation of these settings. This change introduces stricter validation for SSL certificate files (requiring all three or none) and QoS levels, and updates the \MqttClientConnectionFn\ to build connections based on these new settings objects, ensuring that connection options like timeouts, keep-alive intervals, and automatic reconnection are applied consistently for both source and sink operations.

kafka-connect-mqtt/src/main/scala/io/lenses/streamreactor/connect/mqtt/config · high confidence

MongoDB connector configuration and conversion logic restructured

The MongoDB connector's configuration and data handling have been refactored to improve type safety and maintainability. The configuration definition is now centralized in a new \MongoConfig\ object using Kafka Connect's \ConfigDef\, with constants moved to \MongoConfigConstants\. A new \MongoSettings\ case class replaces the previous configuration model, providing a strongly-typed representation of connection details, KCQL mappings, error policies, and SSL settings. Additionally, the \SinkRecordConverter\ has been rewritten to handle Kafka Connect schema types (such as Date, Time, Timestamp, Decimal, and complex Maps/Structs) more robustly, ensuring accurate conversion to MongoDB BSON documents.

kafka-connect-mongodb/src/main/scala/io/lenses/streamreactor/connect/mongodb/config · high confidence

New Avro and JSON sink data converters with enhanced type handling

The cloud sink connectors now use new \ToAvroDataConverter\ and \ToJsonDataConverter\ implementations in the \kafka-connect-cloud-common\ module. The Avro converter explicitly handles logical types (Date, Time, Timestamp, Decimal) and Connect Union structs, ensuring correct schema mapping and default value retrieval. The JSON converter improves serialization by using Jackson for POJO arrays and properly delegating Structs to the Kafka Connect JsonConverter, resolving issues with mixed-type maps and SMT-returned POJOs.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/conversion · high confidence

New cloud connector configuration traits and conversion interface

The connector now uses a new set of Scala traits to define configuration structures for cloud sources and sinks. CloudSinkConfig exposes settings for bucket options, compression, retry policies, error handling, schema change detection, and null-value skipping. CloudSourceConfig adds configuration for partition searching, file extension filtering, backoff settings, watermark header writing, and late-arrival intervals. A new PropsToConfigConverter trait provides a standardized interface for parsing connector properties into these configuration objects, enabling more modular and extensible configuration handling across different cloud storage implementations.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/config/traits · high confidence

New error handling framework with integrity-aware retry policies

The connector now uses a new error handling architecture in \kafka-connect-common\ that introduces \FatalConnectException\ and \RetriableIntegrityException\ to prevent data loss. \FatalConnectException\ signals unrecoverable errors that all policies (NOOP, THROW, RETRY) must fail-fast on, while \RetriableIntegrityException\ allows safe retries for transient integrity issues (like lock timeouts) under the RETRY policy, but still fails fast under NOOP and THROW to avoid silent data loss. The \ErrorHandler\ and \ErrorPolicy\ traits manage retry counts and delegate to specific policy implementations, ensuring consistent behavior across error scenarios.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/common/errors · high confidence

New local output stream implementation with 64-bit pointer tracking

The cloud sink now uses a new \BuildLocalOutputStream\ class in the common stream module to handle local file buffering. This implementation tracks the written byte count using a 64-bit \Long\ pointer, preventing overflow issues for files larger than 2 GiB that could occur with 32-bit integer tracking. The class implements the \CloudOutputStream\ trait, providing a standardized interface for completing writes and reporting the current position.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/stream · high confidence

Refactored cloud connector configuration and storage settings into a shared common module

The configuration logic for cloud connectors has been reorganized into the \kafka-connect-cloud-common\ module to enable reuse across multiple connectors. This change introduces a new \CloudConfigDef\ class that standardizes configuration parsing and applies a lowercase-key processor to ensure case-insensitive property handling. It also centralizes \DataStorageSettings\ to manage envelope and field storage options, \IndexSettings\ to configure exactly-once semantics (including the new \Batch\ commit mode alongside the existing \Granular\ mode) and garbage collection parameters, and \FormatSelection\ to handle format-specific reader construction. Additionally, new utilities like \ConnectorTaskId\ and \TaskDistributor\ are provided to manage task indexing and distribution, while \ConsumerGroupsWriter\ is added to support storing consumer offsets in cloud storage.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/config · high confidence

Refactored cloud source connector file reading into a common reusable module

The file reading logic for Data Lake source connectors has been extracted from connector-specific implementations into a shared \kafka-connect-cloud-common\ module. This change introduces \PartitionDiscovery\ to handle dynamic directory scanning and \ReaderManager\ to orchestrate the lifecycle of file readers, allowing connectors to reuse the same robust state management, polling loops, and post-processing hooks. Users benefit from consistent behavior across different cloud providers (such as AWS S3 and GCP) and improved reliability in handling file discovery and record retrieval.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source/reader · high confidence

Refactored cloud source file listing into a common, reusable module

The file-listing logic for cloud sources has been moved into the \kafka-connect-cloud-common\ module to enable reuse across different cloud providers. This change introduces a \BatchLister\ abstraction with two implementations: \DefaultOrderingBatchLister\ (using native UTF-8 binary ordering) and \DateOrderingBatchLister\ (sorting by file modification date in memory). It also adds a \CloudSourceFileQueue\ to manage the consumption of file batches and an \ExponentialBackoffSourceFileQueue\ to handle polling delays when no files are found, ensuring consistent and resilient file ingestion.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source/files · high confidence

Standardized connector configuration traits and constants

The \kafka-connect-common\ module now provides a unified set of configuration traits and constants for all connectors. A new \TraitConfigConst\ object centralizes property suffixes (such as \kcql\, \error.policy\, \retry.interval\, and SSL settings), while the \BaseConfig\ class and \BaseSettings\ trait establish a consistent foundation for connector configuration. Specific traits like \ErrorPolicySettings\, \RetryConfigSettings\, \KcqlSettings\, \SSLSettings\, and \UserSettings\ are introduced to standardize how common connector options are defined, validated, and accessed, ensuring a consistent configuration experience across the Stream Reactor suite.

kafka-connect-common/src/main/scala/io/lenses/streamreactor/common/config/base · high confidence

Unified cloud source configuration and enhanced file processing controls

The cloud source connectors now use a common configuration module that introduces a unified \source.partition.search.depth\ setting to replace the cloud-specific \recurse.levels\ parameter, ensuring consistent directory traversal across S3, GCS, and Azure. Users can now filter source files by extension using \source.extension.includes\ and \source.extension.excludes\, and control file ordering via a new \ordering.type\ setting (AlphaNumeric or LastModified). Additionally, the post-processing actions (Move and Delete) now support a \retain\ flag to preserve directory structures after processing, and a \processLateArrival\ flag for the Move action to handle late-arriving data. The connectors also support writing watermark information to Kafka record headers via \source.write.watermark.headers\ and allow configuring late arrival check intervals.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/source/config · high confidence

Fixes

Fixes incorrect schema promotion in AttachLatestSchemaOptimizer

The AttachLatestSchemaOptimizer now correctly handles union type conversions during schema evolution. Previously, when promoting a simple type (like an enum) to a union, the logic could incorrectly select the first matching branch based on type alone, ignoring name specificity. This fix ensures that the most specific matching branch is selected, preventing data adaptation errors and ensuring records are written with the correct schema structure.

kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/optimization · high confidence

Test coverage

Add functional tests for S3 connector compression codecs; Add integration tests for Elastic 7 bulk error handling and writer selection; Added HTTP vs GCS sink throughput benchmarks with mocked egress; Added Logback configuration for tests; Added byte-array test fixtures for integration testing; Added configuration validation tests for MongoDB and MQTT connectors; Added configuration validation tests for the OpenSearch sink connector; Added end-to-end test for the MQTT sink connector; Added functional test for Elastic 7 sink connector; Added functional tests for Cassandra source and sink connectors; Added functional tests for the HTTP Sink connector; Added integration test infrastructure for GCP Storage sink; Added integration test suite for GCP Storage Sink; Added integration test utilities for GCP Storage connector; Added integration tests for Azure Data Lake sink task; Added integration tests for GCP Storage Sink schema evolution handling; Added integration tests for GCP storage directory listing and path existence checks; Added integration tests for JMS source and sink connectors; Added integration tests for the OpenSearch sink connector; Added test coverage for source data converters; Added test coverage for the Elastic 7 sink connector; Added test resources for JMS sink connector validation; Added test resources for embedded Cassandra and transaction validation; Added test utilities and configuration validation tests; Added test utilities for Cyclops Either and Option types; Added tests for AWS S3 directory listing and storage interface operations; Added tests for Azure Datalake location validation; Added tests for Cassandra connector configuration, settings, and sink behavior; Added tests for Cyclops to Scala Either conversion; Added tests for HTTP sink authentication and metrics; Added tests for HTTP sink metrics and reset logic; Added tests for HTTP sink record batching logic; Added tests for KcqlWithFieldsSettings upsert key extraction; Added tests for MongoDB SinkRecordConverter date handling; Added tests for S3 configuration processors; Added tests for S3 source configuration and delete mode settings; Added tests for S3 source configuration parsing; Added tests for SQL transformation and struct handling in Kafka Connect SQL common; Added tests for batch policy conditions and logging behavior; Added tests for cloud sink extractors and path splitting; Added tests for error policy handling of fatal and integrity exceptions; Added tests for schema conversion and projection utilities; Added tests for wallclock SMT partition key validation; Added unit tests for AvroSerializer; Added unit tests for Azure CosmosDB connector configuration and client initialization; Added unit tests for Azure Data Lake storage interface and page iteration; Added unit tests for Azure Service Bus KCQL configuration mapping; Added unit tests for Azure Service Bus source and sink record mappers; Added unit tests for Elastic 6 sink connector components; Added unit tests for FTP file listing and filtering logic; Added unit tests for GCP Storage connector configuration and storage interface; Added unit tests for GCP connector authentication and configuration; Added unit tests for HTTP sink template rendering; Added unit tests for HTTP sink template substitution logic; Added unit tests for JMS connector configuration, converters, and message handling; Added unit tests for JSON conversion and MongoDB sink record processing; Added unit tests for KCQL parsing logic; Added unit tests for Kafka Connect common utilities and configuration; Added unit tests for KcqlProperties configuration parsing; Added unit tests for MQTT source connector topic mapping and shared subscription replication; Added unit tests for MongoDB sink key extraction and connector task configuration; Added unit tests for MqttSinkConnector and MqttWriter; Added unit tests for OpenSearch authentication components; Added unit tests for ReaderManagerBuilder; Added unit tests for RecordConverter; Added unit tests for S3 connector configuration and connector logic; Added unit tests for S3 model validation and partition extraction; Added unit tests for S3 sink configuration and validation logic; Added unit tests for S3 source reader components; Added unit tests for cloud-common configuration and data handling components; Added unit tests for concurrent execution, SSL configuration, and offset handling; Added unit tests for text input readers and utilities; Added unit tests for the Azure Service Bus sink connector components; Added unit tests for the GCP Pub/Sub Source Connector; Added unit tests for the HTTP Sink connector; Added unit tests for the KCQL target type parser; Added unit tests for the new OpenSearch sink connector; Added unit tests for the sink-reporting module; Cassandra integration tests migrated to Scala; FTP source integration tests migrated to Scala; Functional tests for Elastic6 and MongoDB connectors migrated to Scala; Integration tests for HTTP Sink and request sender; Integration tests for MongoDB sink migrated to Scala with new test fixtures; Migrate MQTT connector integration tests from Java to Scala; New integration tests for Cassandra sink connector features.

Dependencies

Build system migration to SBT 1.11.2 with consolidated dependency management

The project has migrated its build infrastructure to SBT version 1.11.2, introducing a centralized \Dependencies.scala\ file that explicitly defines versions for all libraries and build plugins. This change consolidates dependency management, setting Scala 2.13.16, Kafka 4.1.0, and AWS SDK 2.29.52, while updating key plugins such as \sbt-assembly\ to 2.3.1, \sbt-scoverage\ to 2.3.1, and \sbt-pack\ to 0.20. The build now enforces Java 11 compatibility and includes new packaging plugins for license headers and assembly configuration.

project · high confidence

Introduce Gradle build system for Java connectors

The Java connectors module now uses Gradle for building, replacing the previous build configuration. This introduces the Shadow plugin (v8.3.6) for creating fat JARs and Spotless (v7.0.4) for code formatting. The build targets Java 17 and includes dependencies for Azure Service Bus (azure-core 1.55.4, azure-messaging-servicebus 7.17.12), GCP (google-cloud libraries-bom 26.38.0), and Kafka 4.1.0.

(dependencies) · high confidence

Housekeeping

Documentation of Exactly Once Offset and Locking Logic

Added a UML sequence diagram to the \kafka-connect-cloud-common\ module that visualizes the task lifecycle for S3 and GCP storage sinks. The diagram details the initialization, lockfile management, record processing, and offset validation steps required to ensure exactly-once delivery semantics, specifically illustrating how pending operations are tracked and how offset consistency is verified before finalizing uploads.

kafka-connect-cloud-common/src/uml · 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

Score

  • CAI 63 → 63 (-0.2)
  • Rubric changed (rubric-2026.09.8 → rubric-2026.09.16) — scores are not directly comparable.

Lenses

  • Code Health 88 → 88 (-0.1)
  • Architecture 98 → 82 (-16.4)
  • Maturity 71 → 69 (-1.5)
  • Readiness 54 → 55 (+1.1)
  • Security 59 → 59 (+0.0)
  • Domain Modelling 100 → 100 (+0.0)
  • Performance 100 (new)

Resolved (4)

  • Documentation: no architecture or design documentation (README.md)
  • Hotspot: kafka-connect-azure-datalake/src/main/scala/io/lenses/streamreactor/connect/datalake/storage/DatalakeStorageInterface.scala (kafka-connect-azure-datalake/src/main/scala/io/lenses/streamreactor/connect/datalake/storage/DatalakeStorageInterface.scala)
  • Off-boarding risk: anonymized user #1
  • Off-boarding risk: anonymized user #2

New (29)

  • Duplicated block (10 lines × 2) (java-connectors/kafka-connect-query-language/src/main/java/io/lenses/kcql/Kcql.java)
  • Duplicated block (11–12 lines × 3) (java-connectors/kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery/MergeQueries.java)
  • Duplicated block (12–14 lines × 3) (java-connectors/kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery/MergeQueries.java)
  • Duplicated block (13 lines × 2) (java-connectors/kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery/convert/logicaltype/DebeziumLogicalConverters.java)
  • Duplicated block (14–15 lines × 2) (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/writer/Writer.scala)
  • Duplicated block (17 lines × 2) (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/seek/IndexManagerV2.scala)
  • Duplicated block (22–23 lines × 2) (java-connectors/kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery/MergeQueries.java)
  • Duplicated block (5 lines × 2) (java-connectors/kafka-connect-common/src/main/java/io/lenses/streamreactor/common/validation/validators/PatternMatchingSourceNameValidator.java)
  • Duplicated block (7 lines × 3) (java-connectors/kafka-connect-azure-eventhubs/src/main/java/io/lenses/streamreactor/connect/azure/eventhubs/util/KcqlConfigTopicMapper.java)
  • Duplicated block (7–8 lines × 2) (java-connectors/kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery/SchemaManager.java)
  • Duplicated block (8 lines × 2) (java-connectors/kafka-connect-azure-servicebus/src/main/java/io/lenses/streamreactor/connect/azure/servicebus/sink/AzureServiceBusSinkConnector.java)
  • Duplicated block (8–10 lines × 2) (java-connectors/kafka-connect-azure-eventhubs/src/main/java/io/lenses/streamreactor/connect/azure/eventhubs/util/KcqlConfigTopicMapper.java)
  • Duplicated block (8–9 lines × 2) (java-connectors/kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery/convert/BigQuerySchemaConverter.java)
  • Hotspot: kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/WriterManagerCreator.scala (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/WriterManagerCreator.scala)
  • Hotspot: kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/writer/Writer.scala (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/writer/Writer.scala)
  • Hotspot: kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/writer/WriterManager.scala (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/writer/WriterManager.scala)
  • IndexManagerV2.snapshotLegacyLocks (cognitive 28) (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/seek/IndexManagerV2.scala)
  • IndexManagerV2.startExecutors (cognitive 31) (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/seek/IndexManagerV2.scala)
  • IndexManagerV2.sweepOrphansUnder (cognitive 17) (kafka-connect-cloud-common/src/main/scala/io/lenses/streamreactor/connect/cloud/common/sink/seek/IndexManagerV2.scala)
  • No ADRs found
  • …and 9 more

Changes since last survey

  • 6 commits — 5 feature/other, 1 fixes

By area

  • (repo) — 3 commits
  • kafka-connect-cloud-common/src — 2 commits
  • java-connectors/build.gradle — 1 commit

Notable commits

  • fix: Merge pull request #423 from lensesio-dev/fix/wallclock-partitionby-requires-batch-commit-mode
  • change: Merge pull request #419 from lensesio-dev/feat/partition-batch-commit-mode
  • change: Merge pull request #420 from lensesio-dev/chore/OPS-2652-drop-never-loaded-libraries
  • change: Reject non-deterministic wallclock PARTITIONBY keys under granular commit mode
  • change: chore(build): stop packaging libraries that can never load
  • change: fail closed on unreadable master lock to prevent orphan-sweep from deleting in-flight Copy sources

Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.

Survey your own repository

lensesio/stream-reactor 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 28 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 8d37476a9c7e5ae7d49d77b6bf3cd72e9a6cb568 — the exact code this score is about.
  • Scored under rubric-2026.09.16 — the same rubric and the same method as every other entry in this index.
  • Measured by watchdog.canine.dev using codehealth-analyzer preprod-2d9048c36d26.