kafka-aggregator
A Kafka aggregator based on the Faust Python Stream Processing library.
kafka-aggregator development is based on the Safir application template.
Overview
kafka-aggregator uses Faust’s windowing feature to aggregate a stream of messages from Kafka.
kafka-aggregator implements a Faust agent, a “stream processor”, that adds messages from a source topic into a Faust table. The table is configured as a tumbling window with a size, representing the window duration (time interval) and an expiration time, which specifies the duration for which the data allocated to each window will be stored. Every time a window expires, a callback function is called to aggregate the messages allocated to that window. The size of the window controls the frequency of the aggregated stream.
kafka-aggregator uses faust-avro to add Avro serialization and Schema Registry support to Faust. faust-avro can parse Faust models into Avro Schemas.
See the docs for more information.
Change log
0.2.0 (2020-08-14)
Add first and third quartiles (q1 and q3) to the list of summary statistics computed by the aggregator.
Ability to configure the list of summary statistics to be computed.
Pinned top-level requeriments.
Add Kafka Connect to the docker-compose setup.
Use only one Schema Registry by default to simplify local execution.
First release to PyPI.
0.1.0 (2020-07-13)
Initial release of kafka-aggregator with the following features:
Use Faust windowing feature to aggregate a stream of messages.
Use Faust-avro to add Avro serialization and Schema Registry support to Faust.
Support to an internal Schema Registry to store schemas for the aggreated topics (optional).
Create aggregation topic schemas from the source topic schemas and from the list of summary statistics to be computed.
Ability to create Faust records dynamically from aggregation topic schemas.
Ability to auto-generate code for the Faust agents (stream processors).
Compute summary statistics for numeric fields: min(), mean(), median(), stdev(), max().
Add example module to initialize a number of source topics in kafka, control the number of fields in each topic, and produce messages for those topics at a given frequency.
Use Kafdrop to inspect messages from source and aggregated topics.
Add kafka-aggregator documentation site.
Metadata
Release files for kafka-aggregator 0.2.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| kafka-aggregator-0.2.0.tar.gz | 78.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| kafka_aggregator-0.2.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 96.8 kB
Release files / kafka-aggregator-0.2.0.tar.gz
| Download URL | kafka-aggregator-0.2.0.tar.gz |
|---|---|
| Size | 78.7 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
0189b674d6468360652c2efbff322ee18987cdfadb21f6d41ce843609b95dbb4
|
|
BLAKE2b-256 checksum How to use checksums |
03aed4ea6da57518a3760900f397191613e66b23166044a3d5017a96172452cd
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/3.2.0 pkginfo/1.5.0.1 requests/2.24.0 setuptools/49.3.1 requests-toolbelt/0.9.1 tqdm/4.48.2 CPython/3.8.5
|
Release files / kafka_aggregator-0.2.0-py3-none-any.whl
| Download URL | kafka_aggregator-0.2.0-py3-none-any.whl |
|---|---|
| Size | 18.1 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
24e0ab5853788d22d9654be3ca4e4e0793b11f5b37421bde9a77a2b53d9c8e81
|
|
BLAKE2b-256 checksum How to use checksums |
8f4be670e57e52d20e5edab7453fe4fa9cf2fc9e767165f1dd686173f2257dca
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/3.2.0 pkginfo/1.5.0.1 requests/2.24.0 setuptools/49.3.1 requests-toolbelt/0.9.1 tqdm/4.48.2 CPython/3.8.5
|