quixio/quix-streams
63.1
Adequate · 22 September 2026
84k
lines of production code
Python
primary language
7
measurements over time
What this system is
Quix Streams is a Python library for building declarative, stateful stream processing applications on top of Apache Kafka. It provides a Pandas-like API for defining data transformation pipelines, supporting complex operations such as time-based joins, windowed aggregations, and lookups against external configuration services. The system manages the underlying Kafka consumer/producer lifecycle, including exactly-once semantics, checkpointing, and robust state recovery using RocksDB. It also includes a comprehensive ecosystem of connectors for ingesting data from sources like MySQL CDC and Kinesis, and exporting processed results to sinks such as BigQuery, InfluxDB, and various cloud storage services.
How it got here
2023 — quixstreams v3.26.0 release
20 changes.
This period focused on the release of quixstreams v3.26.0, introducing a new core streaming engine with a declarative StreamingDataFrame API and a unified serialization framework for Kafka messages. The update also established robust state management with TTL migration, native Quix platform integration, and comprehensive test coverage for all new subsystems.
2024 — Windowing, Sinks, and Sources expansion
31 changes.
This period focused on expanding the platform's data processing capabilities by introducing comprehensive windowing aggregations with RocksDB-backed state storage. It also established a robust, extensible API for data ingestion (Sources) and output (Sinks), adding support for various cloud storage systems, databases, and message brokers. Concurrently, the core application architecture was refactored to include centralized processing context, improved checkpointing, and enhanced topic management.
2025–2026 — streaming joins and connector expansion
12 changes.
This period focused on expanding the framework's data integration capabilities by introducing new connectors for InfluxDB 3.0, TDengine, and MySQL CDC, alongside a modular refactoring of the file source. Significant engineering effort was directed toward enhancing streaming data processing through the implementation of AsOfJoin, IntervalJoin, and configurable lookup buffering mechanisms. These features were supported by the development of an internal buffered consumer for ordered message processing and comprehensive test coverage across the new components.
Features
Add Amazon Kinesis source connector
Users can now consume data from an Amazon Kinesis stream and forward it to a Kafka topic using the new \KinesisSource\. This feature supports configurable AWS credentials (via parameters or environment variables), automatic offset reset strategies (earliest/latest), batched round-robin consumption with configurable record limits, and checkpointing for at-least-once delivery guarantees. It also includes optional callbacks for client connection success and failure events.
quixstreams/sources/community/kinesis · high confidence
Add AsOfJoin and IntervalJoin capabilities to the StreamingDataFrame
This change introduces two new join types for the StreamingDataFrame: AsOfJoin, which matches records based on the latest available value at a specific timestamp, and IntervalJoin, which matches records falling within specified backward and forward time windows. The implementation includes a base Join class handling state registration and merge strategies (keep-left, keep-right, or raise on conflict), along with utility functions for merging dictionaries. Users can now perform time-based joins on streaming data using these new classes exposed via the quixstreams.dataframe.joins module.
quixstreams/dataframe/joins · high confidence
Add Google Cloud Pub/Sub source connector
Users can now ingest data from Google Cloud Pub/Sub topics into Kafka using the new PubSubSource. This feature introduces a consumer that supports optional service account authentication, automatic subscription creation, message ordering, and configurable commit batching. It also exposes optional callbacks to monitor client connection success or failure.
quixstreams/sources/community/pubsub · high confidence
Add community sinks for BigQuery, Elasticsearch, Iceberg, InfluxDB, Kafka, Kinesis, MongoDB, and MQTT
New sink connectors are now available in the \quixstreams.sinks.community\ module, allowing users to export processed data to a variety of external systems. This release introduces \BigQuerySink\ for Google Cloud BigQuery, \ElasticsearchSink\ for Elasticsearch, \IcebergSink\ for Apache Iceberg tables (currently supporting AWS Glue), \InfluxDB1Sink\ for InfluxDB v1, \KafkaReplicatorSink\ for replicating data to external Kafka clusters, \KinesisSink\ for AWS Kinesis, \MongoDBSink\ for MongoDB, and \MQTTSink\ for MQTT brokers. Each connector handles its specific authentication, batching, and serialization requirements, expanding the platform's data integration capabilities.
quixstreams/sinks/community · high confidence
Add utility modules for JSON, serialization, settings, and table printing
This change introduces a new \quixstreams/utils\ package containing several helper modules. It adds \json.py\ for high-performance JSON serialization using \orjson\, \pickle.py\ for efficient object copying via \pickle\_copier\, and \settings.py\ providing a \BaseSettings\ class with Pydantic-based configuration handling. Additionally, it includes \printing.py\ which implements a \Printer\ and \Table\ class for displaying data in console applications (supporting both interactive and non-interactive modes), \dicts.py\ for recursive dictionary flattening, and \stream\_id.py\ for generating deterministic stream identifiers.
quixstreams/utils · high confidence
Initial Conda packaging support for Quix Streams
Quix Streams is now available as a Conda package, allowing users to install and manage the library via the Conda environment manager. This change introduces the necessary build infrastructure, including a \meta.yaml\ specification file that defines dependencies (such as Pydantic, httpx, and various cloud SDKs) and a \post-link.sh\ script to handle packages not available in standard Conda channels. It also provides a release script for publishing to Anaconda.org and a validation tool to ensure consistency between PyPI and Conda requirements.
conda · high confidence
Introduce KafkaReplicatorSource and QuixEnvironmentSource for topic replication
Added new source implementations that replicate messages from an external Kafka broker or a Quix Cloud environment into the application's Kafka cluster. The \KafkaReplicatorSource\ allows users to connect to a remote broker and stream topics, while \QuixEnvironmentSource\ simplifies this for Quix Cloud workspaces by handling authentication and topic resolution automatically. Both sources support configurable deserializers, error callbacks, and optional connection success/failure hooks, enabling safe data copying for development, testing, or cross-environment synchronization.
quixstreams/sources/core/kafka · high confidence
Introduce ProcessingContext for centralized application state management
A new \ProcessingContext\ class has been added to the \quixstreams.processing\ module to serve as a central hub for sharing processing-related objects between the \Application\ and \StreamingDataFrame\ instances. This context manages critical components including the internal consumer and producer, state store manager, sink manager, and dataframe registry, while also handling checkpoint initialization, offset storage, and commit logic. It ensures proper lifecycle management by starting sinks on entry and handling transaction abortion or printer cleanup on exit, providing a unified interface for core processing state.
quixstreams/processing · high confidence
Introduce Quix Configuration Service lookup for streaming enrichment
Adds a new \QuixConfigurationService\ lookup join that enriches streaming records with versioned configuration data from a Kafka topic. The implementation supports extracting JSON fields via JSONPath and binary (bytes) fields, handles configuration updates and deletions, and includes an LRU cache for configuration version data. It also supports optional SDK token authentication for secure access to configuration content.
_quixstreams/dataframe/joins/lookups/quix\_configuration\service · high confidence
Introduce Quix platform integration module with API client and topic management
Added the \quixstreams.platforms\ package to provide native integration with the Quix Cloud platform. This includes a \QuixPortalApiService\ for communicating with the Quix API, a \QuixKafkaConfigsBuilder\ to automatically retrieve broker connection details and TLS certificates, and a \QuixTopicManager\ that handles topic creation and configuration specific to Quix workspaces. The module also introduces environment variable handling for Quix-specific settings (such as workspace ID and deployment state) and dedicated exception types for platform-related errors.
quixstreams/platforms · high confidence
Introduce RocksDB-backed windowed state store
Added a new windowed state implementation using RocksDB, including the \WindowedRocksDBStore\, \WindowedRocksDBStorePartition\, and \WindowedRocksDBPartitionTransaction\ classes. This change introduces specific column families for managing window metadata (such as expiration and deletion indexes) and provides a \WindowedTransactionState\ interface that allows users to get, update, and expire windows, as well as manage collection-type aggregations within the windowed state.
quixstreams/state/rocksdb/windowed · high confidence
Introduce Sources API with CSV and Kafka Replicator connectors
The \quixstreams.sources\ module is now available, exposing a new public API for data ingestion. Users can now utilize \CSVSource\ for reading CSV files, as well as \KafkaReplicatorSource\ and \QuixEnvironmentSource\ for Kafka-based data streams. The module also provides base classes like \BaseSource\ and \StatefulSource\, along with \SourceManager\ and connection callbacks (\ClientConnectSuccessCallback\, \ClientConnectFailureCallback\) to manage source lifecycle and connectivity.
quixstreams/sources · high confidence
Introduce configurable lookup join buffering for streaming dataframes
Added a new \LookupBuffer\ mechanism to the \join\_lookup\ API that allows incoming records to be held in a stateful buffer when their external configuration is not yet available, rather than being immediately enriched with defaults or dropped. Users can now specify a \grace\_ms\ window, a maximum buffer size per key, and a timeout strategy (\emit\ with defaults or \drop\) to manage latency and data completeness. The buffer persists withheld records in a changelog-backed state store (RocksDB), ensuring state recovery across restarts, and includes a clock-driven tick to release timed-out records even on idle partitions.
quixstreams/dataframe/joins/lookups · high confidence
Introduce core sink implementations and infrastructure
Adds the foundational sink components for the platform: a unified BlobStorageClient for accessing Azure, AWS S3, GCP, and local storage via quixportal; a REST Catalog client for table registration; and specific sinks including InfluxDB3Sink (with configurable time precision, SSL verification, and integer-to-float conversion), QuixTSDataLakeSink (writing Hive-partitioned Parquet files with virtual partitions, column statistics, and stream-timeout detection), CSVSink for local debugging, and ListSink for interactive inspection.
quixstreams/sinks/core · high confidence
Introduce in-memory state store and refactor RocksDB state backend
Added a new in-memory state store (MemoryStore) that mirrors the semantics and TTL configuration of the RocksDB backend, providing a lightweight alternative for development and testing. Refactored the RocksDB state backend by extracting a base store implementation, introducing a shared block cache to optimize memory usage across partitions, and adding an open-deadline mechanism to prevent lock contention during Kafka rebalances. The update also includes comprehensive TTL migration logic to handle legacy stores, ensuring data integrity during upgrades.
quixstreams/state/rocksdb · high confidence
Introduce windowed aggregations and time/count-based window definitions
The \quixstreams/dataframe/windows\ module now provides a complete set of windowing capabilities, including tumbling, hopping, and sliding windows for both time and count triggers. Users can apply built-in aggregations such as \sum\, \count\, \mean\, \max\, \min\, and \reduce\, or collect all values within a window period. The implementation supports both \final\ (closed window) and \current\ (streaming) output modes, with configurable grace periods and closing strategies (key-based or partition-based) to handle late data and window expiration.
quixstreams/dataframe/windows · high confidence
Introduction of core data models for Kafka message processing
This change introduces the foundational data structures for the quixstreams library, defining how Kafka messages and their metadata are represented in Python. It adds the \MessageContext\ class to hold immutable Kafka-specific properties like topic, partition, offset, and leader epoch, and the \KafkaMessage\ class to represent raw message components including key, value, headers, and timestamp. The \Row\ class is introduced as the primary unit of data flowing through the streaming pipeline, combining the message value with its context and metadata. Additionally, \MessageTimestamp\ and \TimestampType\ are added to handle Kafka message timestamps with support for creation time and log append time, while \types.py\ defines type aliases and protocols to ensure type safety and compatibility with the underlying \confluent\_kafka\ library.
quixstreams/models · high confidence
Introduction of the StreamingDataFrame API for declarative stream processing
The \quixstreams.dataframe\ module has been introduced, providing the \StreamingDataFrame\ and \StreamingSeries\ classes that enable declarative, Pandas-like ETL pipelines for Kafka streams. This new API allows users to build processing topologies by chaining operations such as filtering, applying transformations, and joining data, while the underlying \DataFrameRegistry\ manages topic subscriptions, stream IDs, and timestamp alignment requirements for operations like concatenation and joins.
quixstreams/dataframe · high confidence
New FileSink with local, S3, and Azure storage backends
A new FileSink component is introduced, allowing data batches to be written to files in configurable formats (JSON Lines or Parquet) across multiple storage destinations. The implementation includes LocalFileSink for writing to the local filesystem (with optional append mode), S3FileSink for Amazon S3 buckets, and AzureFileSink for Azure Blob Storage containers. Users can configure serialization formats, directory structures, and authentication methods for each backend.
quixstreams/sinks/community/file · high confidence
New FileSource connectors for local, S3, and Azure Filestore
This change introduces a new file-based source capability within the \quixstreams.sources\ module, allowing users to ingest records from files stored in local filesystems, AWS S3 buckets, and Azure Filestore containers. The implementation provides \LocalFileSource\, \S3FileSource\, and \AzureFileSource\ classes that recursively iterate through file paths, parse records in JSON or Parquet formats (with optional GZip decompression), and produce them to Kafka topics. Users can configure replay speeds, partition folder alignment, and custom key/value/timestamp setters to control how file data is streamed into their Kafka pipelines.
quixstreams/sources/community/file · high confidence
New InfluxDB 3.0 data source connector
Added a new \InfluxDB3Source\ component in the \quixstreams/sources/community/influxdb3\ module, enabling users to extract data from InfluxDB 3.0 instances. This connector supports querying specific or all measurements using time-windowed tumbling intervals, allows for custom SQL queries, and includes built-in retry logic with exponential backoff for robust data ingestion.
quixstreams/sources/community/influxdb3 · high confidence
New MySQL CDC Lite source for streaming binlog changes
A new \MySqlCdcLiteSource\ component has been added to the community sources, enabling streaming of MySQL binary log (binlog) changes. This source connects to a MySQL server to capture row-level events (inserts, updates, deletes) and publishes them as change records. It handles data encoding for various MySQL types, including JSON, BINARY, and TIME with fractional seconds, and manages connection state and retries. This feature is available as a streaming-only source.
_quixstreams/sources/community/mysql\_cdc\lite · high confidence
New Sinks API with connection lifecycle callbacks
The quixstreams.sinks module has been reorganized to expose a new public API, making BaseSink, BatchingSink, SinkManager, and related classes available for import. This update introduces optional ClientConnectSuccessCallback and ClientConnectFailureCallback types, allowing users to handle connection lifecycle events for sinks.
quixstreams/sinks · high confidence
New TDengine sink for streaming data
Added a new TDengineSink component that allows users to stream processed data to a TDengine database. The sink batches records, converts them to the InfluxDB line protocol, and supports configurable fields, tags, time precision, and authentication methods (token or username/password).
quixstreams/sinks/community/tdengine · high confidence
New community sources for MQTT and Pandas DataFrames
Users can now ingest data from MQTT brokers and local Pandas DataFrames using the new MQTTSource and PandasDataFrameSource components in the community sources module. The MQTTSource connects to brokers (supporting versions 3.1, 3.1.1, and 5) with configurable TLS, authentication, and QoS, and exposes optional callbacks for connection success and failure events. The PandasDataFrameSource reads rows from a DataFrame and produces them as JSON messages to Kafka, allowing users to specify key and timestamp columns, apply a row-by-row delay for simulation, and choose whether to include metadata in the message values.
quixstreams/sources/community · high confidence
New internal buffered consumer for ordered message processing
Added a new \InternalConsumer\ class in the \quixstreams/internal\_consumer\ module that wraps the underlying Kafka consumer with a buffering layer. This consumer maintains an in-memory buffer of messages per partition to ensure strict ordering when consuming from multiple topics, preventing stalls after restarts by tracking high watermarks and consumer positions. Users relying on internal streaming mechanics will now benefit from this buffered approach, which manages backpressure and idleness detection automatically.
_quixstreams/internal\consumer · high confidence
New serialization framework with Avro, Protobuf, and JSON Schema support
The \quixstreams.models.serializers\ package has been introduced, providing a unified serialization and deserialization (SerDes) system for Kafka messages. This update adds native support for Avro, Protobuf, and JSON Schema formats, allowing users to serialize data using specific schemas and integrate with Confluent Schema Registry for schema management and validation. The package includes base classes for custom serializers, built-in handlers for simple types (strings, integers, doubles, bytes), and specialized serializers for Quix-specific data models (Timeseries and Events). Users can now configure serialization behavior, including schema validation, deterministic Protobuf encoding, and Schema Registry subject naming strategies, directly within the serializer constructors.
quixstreams/models/serializers · high confidence
Behavioural changes
Extracted base state store abstractions and migration delivery logic
The state management layer has been refactored to introduce a new \quixstreams.state.base\ package containing abstract base classes (\Store\, \StorePartition\, \State\, \TransactionState\) and a shared migration delivery loop (\migration\_flush.py\). This change extracts the common implementation details previously embedded in the RocksDB store, providing a standardized interface for state operations and ensuring consistent, crash-safe handling of legacy-to-TTL migration flushes across different storage backends.
quixstreams/state/base · high confidence
Introduce ConnectionConfig and broker-unavailability detection for Kafka consumers and producers
The \quixstreams/kafka\ module now exposes a \ConnectionConfig\ class that centralizes librdkafka connection settings (bootstrap servers, SASL, SSL, OAuth) and handles secret masking and case normalization. Consumers and producers automatically track broker connectivity via librdkafka stats; if all brokers are down, the client records the unavailability time, which enables the application to detect prolonged outages and handle restarts more robustly. Additionally, the consumer defaults to the "range" partition assignment strategy, and the producer includes configurable poll timeouts and retry logic to mitigate \BufferError\ conditions.
quixstreams/kafka · high confidence
Introduce modular file source components for deserialization and fetching
The file source implementation has been refactored into distinct, reusable components within the \quixstreams/sources/community/file/components\ directory. Users can now leverage \FileFetcher\ to handle concurrent file downloads and streaming via a background thread pool, while \raw\_filestream\_deserializer\ provides a standardized way to parse file content into iterables based on specified formats and compression settings. This change separates the concerns of data retrieval and data parsing, making the file source logic more modular and easier to maintain.
quixstreams/sources/community/file/components · high confidence
Introduce new Sinks API with base classes and backpressure handling
The \quixstreams.sinks.base\ module now provides the foundational classes for the new Sinks API, including \BaseSink\ and \BatchingSink\ for implementing data output destinations, \SinkManager\ for lifecycle management, and \SinkBatch\ for accumulating records. This change introduces explicit support for client connection callbacks (\on\_client\_connect\_success\, \on\_client\_connect\_failure\) and a new \SinkBackpressureError\ exception, allowing applications to handle connection failures and destination backpressure by pausing topic partitions.
quixstreams/sinks/base · high confidence
Introduce new Topic management API with admin capabilities and type safety
The \quixstreams.models.topics\ module has been restructured to provide a more robust and type-safe interface for managing Kafka topics. This change introduces a dedicated \TopicAdmin\ class for cluster-level operations such as listing, inspecting, and creating topics, including handling of authorization errors and timeouts. A new \TopicManager\ centralizes topic registration, supporting regular, repartition, and changelog topics with configurable defaults for partitions and replication. The \Topic\ model now explicitly distinguishes between \REGULAR\, \REPARTITION\, and \CHANGELOG\ types via a \TopicType\ enum, and ensures configuration objects are deep-copied to prevent mutation. Additionally, new exception classes like \TopicPermissionError\ and \TopicConfigurationMismatch\ provide clearer error reporting for topic management failures.
quixstreams/models/topics · high confidence
Introduce new checkpointing module with configurable commit behavior
The \quixstreams/checkpointing\ package has been added, introducing a new checkpointing system that supports both time-based (\commit\_interval\) and count-based (\commit\_every\) commit triggers. This change modifies how application state and consumer offsets are persisted, allowing users to configure commit frequency more granularly. The implementation includes new exception types for handling invalid offsets and commit errors, and integrates with the internal consumer/producer and state store managers to manage exactly-once processing semantics and revoke path timeouts.
quixstreams/checkpointing · high confidence
Introduce new core streaming engine with function-based stream processing
The \quixstreams/core\ module now provides a new internal streaming engine that replaces the previous implementation. This change introduces a \Stream\ class that models data processing as a directed acyclic graph of functions, supporting four operation types: Apply (transform values), Update (mutate values in-place), Filter (drop values based on a predicate), and Transform (modify values, keys, and timestamps). The engine includes dedicated function wrappers (\ApplyFunction\, \FilterFunction\, etc.) that handle execution logic, including branching support that automatically copies values when a stream splits into multiple downstream paths. This refactors the core data flow mechanics while maintaining the same high-level API surface for users.
quixstreams/core · high confidence
Introduce new source execution model with multiprocessing and stateful support
Sources are now executed in separate child processes (using the 'spawn' start method) managed by a new SourceManager, which handles lifecycle, signal handling, and inter-process communication. This change introduces a BaseSource/Source interface with a defined setup method for client connection validation, optional success/failure callbacks, and support for StatefulSource implementations that automatically recover state from changelog topics during initialization.
quixstreams/sources/base · high confidence
New state management subsystem with TTL migration and robust recovery
The \quixstreams/state\ package has been introduced to centralize state store management, recovery, and serialization. This change adds a new \StateStoreManager\ that coordinates state stores (including RocksDB, Memory, and Timestamped variants) and handles changelog-based recovery. A key behavioral update is the automatic migration of legacy state stores to a Time-To-Live (TTL) format, ensuring data is preserved during the transition while enabling expiration semantics. The recovery process has been hardened to survive rebalances and handle edge cases like invalid offsets or incomplete migrations, with specific exceptions raised for format incompatibilities or contradictory operational levers to prevent data loss.
quixstreams/state · high confidence
Quixstreams v3.26.0 release with state logging improvements
This release updates the quixstreams library to version 3.26.0. A key behavioral change is the introduction of a dedicated log handler for the 'quixstreams.state' namespace, allowing operators to increase state-related logging verbosity via the QUIXSTREAMS\_STATE\_LOG\_LEVEL environment variable without affecting the rest of the application's logging. This is achieved by attaching a separate handler to the state logger while keeping propagation enabled to ensure errors are not silently dropped.
quixstreams · high confidence
Test coverage
Added comprehensive test suite for state store partitioning, recovery, and TTL migration; Added comprehensive tests for RocksDB legacy-TTL migration and backfill resilience; Added test coverage for DataFrame join operations and lookup buffering; Added test coverage for Quix and Schema Registry serializers; Added test coverage for TopicManager, TopicAdmin, and Topic deserialization; Added test coverage for core sink implementations; Added test infrastructure and fixtures for Kafka, Schema Registry, and application models; Added test package initialization files for processing and sinks modules; Added tests for CSVSource and KafkaReplicatorSource; Added tests for Kafka broker availability, connection configuration, and producer/consumer behavior; Added tests for MQTT and PostgreSQL community sinks; Added tests for MySQL CDC Lite Source; Added tests for Quix platform API, configuration, and topic management; Added tests for base sink components; Added tests for core stream processing functions and stream topology; Added tests for state recovery robustness and edge cases; Added tests for state store migration and TTL behavior; Added tests for the Printer utility class; Added tests for the SourceManager lifecycle and error handling; Added tests for the internal consumer and its buffering logic; Added tests for window aggregations, base window logic, and time-based window types; Added tests for windowed RocksDB state operations; Initial test suite for Quixstreams core components; Initial test suite for StreamingDataFrame and related components.
Dependencies
Migrate to pyproject.toml and update core dependencies
The project has migrated its build configuration from legacy setup files to a modern pyproject.toml, introducing structured optional dependency groups for connectors such as AWS, Azure, BigQuery, Elasticsearch, and MQTT. Core runtime dependencies have been updated, notably upgrading pydantic to support versions up to 2.14 and pydantic-settings up to 2.16, while switching the HTTP client from requests to httpx. Development tooling now relies on pre-commit (\>=3.4,\<4.4) and mypy 1.18.2, with specific type stubs for protobuf and jsonschema also updated.
(dependencies) · high confidence
Housekeeping
Added third-party license files for AWS SDK, Confluent Kafka, FastAvro, Google Cloud, httpx, InfluxDB, and jsonlines
The LICENSES directory now includes explicit license text files for several key dependencies: boto3, confluent-kafka-python, fastavro, google-cloud-bigquery, google-cloud-pubsub, httpx, influxdb3-python, and jsonlines. This ensures that the terms of use for these libraries are readily available for compliance and distribution purposes.
LICENSES · 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 59 → 63 (+4.0)
- Rubric changed (rubric-2026.08.19 → rubric-2026.09.15) — scores are not directly comparable.
Lenses
- Code Health 85 → 89 (+3.7)
- Architecture 100 → 98 (-1.8)
- Maturity 54 → 54 (+0.2)
- Readiness 50 → 53 (+2.6)
- Security 72 → 90 (+18.1)
Resolved (46)
- Application._on_assign (cognitive 28) (quixstreams/app.py)
- Change coupling: exceptions.py ↔ partition.py (quixstreams/state/rocksdb/exceptions.py)
- Coverage not included — suite not readable by the collector
- Dependency hygiene not measured — dependency manifest found but not parsed for hygiene
- Duplicated block (10 lines × 2) (quixstreams/dataframe/windows/definitions.py)
- Duplicated block (10 lines × 2) (quixstreams/state/memory/partition.py)
- Duplicated block (11 lines × 2) (quixstreams/sinks/community/tdengine/sink.py)
- Duplicated block (11 lines × 2) (quixstreams/state/rocksdb/windowed/transaction.py)
- Duplicated block (12 lines × 2) (quixstreams/sinks/core/stream_timeout_tracker.py)
- Duplicated block (13 lines × 2) (quixstreams/dataframe/windows/definitions.py)
- Duplicated block (13 lines × 5) (quixstreams/sinks/community/bigquery.py)
- Duplicated block (14 lines × 2) (quixstreams/kafka/consumer.py)
- Duplicated block (14 lines × 2) (quixstreams/state/memory/partition.py)
- Duplicated block (14 lines × 2) (quixstreams/state/memory/partition.py)
- Duplicated block (15 lines × 2) (quixstreams/state/memory/partition.py)
- Duplicated block (15 lines × 2) (quixstreams/state/memory/partition.py)
- Duplicated block (16 lines × 2) (quixstreams/dataframe/joins/lookups/postgresql.py)
- Duplicated block (16 lines × 2) (quixstreams/kafka/consumer.py)
- Duplicated block (19 lines × 2) (quixstreams/sinks/community/influxdb1.py)
- Duplicated block (20 lines × 2) (quixstreams/dataframe/joins/lookups/postgresql.py)
- …and 26 more
New (140)
- Application._assign_partitions (cognitive 27) (quixstreams/app.py)
- Change coupling: consumer.py ↔ producer.py (quixstreams/kafka/consumer.py)
- Dependency hygiene PARTLY measured — Python dependencies read, no exact pin to grade for currency
- Documentation: no contributor guidance (docs/connectors/contribution-guide.md)
- Documentation: no installation or build instructions (README.md)
- Documentation: no project overview (docs/api-reference/sources.md)
- Duplicated block (10 lines × 3) (quixstreams/sinks/community/influxdb1.py)
- Duplicated block (10–29 lines × 2) (quixstreams/state/memory/partition.py)
- Duplicated block (11 lines × 2) (quixstreams/dataframe/joins/lookups/postgresql.py)
- Duplicated block (11 lines × 2) (quixstreams/dataframe/joins/lookups/postgresql.py)
- Duplicated block (11 lines × 2) (quixstreams/dataframe/windows/definitions.py)
- Duplicated block (11–12 lines × 2) (quixstreams/state/rocksdb/windowed/transaction.py)
- Duplicated block (11–12 lines × 4) (quixstreams/state/memory/partition.py)
- Duplicated block (12 lines × 3) (quixstreams/sinks/community/influxdb1.py)
- Duplicated block (13 lines × 2) (quixstreams/dataframe/windows/definitions.py)
- Duplicated block (13–14 lines × 2) (quixstreams/sinks/core/stream_timeout_tracker.py)
- Duplicated block (14 lines × 2) (quixstreams/kafka/consumer.py)
- Duplicated block (14 lines × 2) (quixstreams/sinks/community/tdengine/sink.py)
- Duplicated block (14 lines × 5) (quixstreams/sinks/community/bigquery.py)
- Duplicated block (16 lines × 2) (quixstreams/state/memory/partition.py)
- …and 120 more
Changes since last survey
- 17 commits — 16 feature/other, 1 fixes
By area
- (root) — 6 commits
- quixstreams/state — 3 commits
- docs/connectors — 2 commits
- docs/api-reference — 1 commit
- quixstreams/init.py — 1 commit
- quixstreams/dataframe — 1 commit
- quixstreams/sources — 1 commit
- tests/requirements.txt — 1 commit
- tests/test_quixstreams — 1 commit
Notable commits
- fix: fix(state): apply configured RocksDB options to existing column families (#1137)
- change: Bump testcontainers[postgres] from 4.13.2 to 4.13.3 (#1066)
- change: Bump types-protobuf from 6.32.1.20250918 to 6.32.1.20251105 (#1065)
- change: Migrate legacy state stores to TTL without data loss (#1134)
- change: New Connector: MySQL CDC Source (streaming only) (#1152)
- change: RecoveryManager: survive a rebalance during changelog recovery without losing messages (sc-74868) (#1146)
- change: State TTL migration: always migrate stamped stores, never double-wrap (sc-74843) (#1145)
- change: Update azure-storage-blob requirement from <12.25,>=12.24.0 to >=12.24.0,<12.28 (#1062)
- change: Update confluent-kafka[avro,json,protobuf,schemaregistry] requirement from <2.12,>=2.8.2 to >=2.8.2,<2.13 (#1064)
- change: Update documentation (#1138)
- change: Update elasticsearch requirement from <9,>=8.17 to >=8.17,<10 (#1063)
- change: chore(deps): support pydantic 2.13 and pydantic-settings 2.15 (#1155)
- change: chore: Version bump to 3.26.0 (#1150)
- change: docs(contributing): require regression tests and docs for user-visible changes (#1154)
- change: feat(sinks): virtual partitions, file statistics and sort column for QuixLake sink (#1143)
- change: feat: Column stats, virtual partitions, sort column and row-group sizing for QuixTSDataLakeSink (#1147)
- change: refactor(lookup): drop runtime data validation from the buffer (#1151)
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
quixio/quix-streams 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 22 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 6c87346c5b1a23b90fbc1c53bb960661bea71467 — 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-821afab8930d.