Kafnus connect container
7.9K
Coverage badge scope: custom Java SMTs (HeaderRouter and MongoNamespacePrefix), merging unit tests and Python E2E tests.
Kafnus Connect is the persistence layer of the Kafnusβ ecosystem β a modern, Kafka-based replacement for Cygnus in FIWARE smart city environments.
It provides ready-to-use Kafka Connect images with custom Single Message Transforms (SMTs) and pre-integrated sink connectors for PostGIS, MongoDB, and HTTP endpoints.
This project is part of FIWAREβ . For more information check the FIWARE Catalogue entry for the Core Context Managementβ .
| :whale: Docker Hubβ |
|---|
Kafnus Connect consumes processed NGSI events from Kafka topics (produced by Kafnus NGSIβ ) and persists them into target datastores or APIs.
Kafka (processed topics)
β
Header Router (datamodels)
β
βΌ
Kafnus Connect (Kafka Connect)
ββ JDBC Sink (PostGIS)
ββ MongoDB Sink
ββ HTTP Sink
The HeaderRouter SMT is responsible for dynamically resolving the destination database schema and table name at runtime based on NGSI headers and a configurable SQL datamodel, removing any SQL layout logic from upstream producers.
Each connector can be independently configured via environment variables or connect-distributed.properties.
Custom SMTs can be chained to transform headers or message formats before persistence.
docker build -t telefonicaiot/kafnus-connect:latest .
docker run -d --name kafnus-connect -e CONNECT_BOOTSTRAP_SERVERS=kafka:9092 -e CONNECT_GROUP_ID=kafnus-connect -e CONNECT_CONFIG_STORAGE_TOPIC=connect-configs -e CONNECT_OFFSET_STORAGE_TOPIC=connect-offsets -e CONNECT_STATUS_STORAGE_TOPIC=connect-status telefonicaiot/kafnus-connect:latest
For complete examples, see the
tests_end2endβ folder in the main Kafnus repository.
Integration and end-to-end testing are performed from the Kafnus NGSIβ repository, where complete data flow scenarios are executed using Testcontainers.
This repository also includes his own python tests (similar to Kafnus tests) and unit tests for the custom Java SMTs in src/kafnus-connect-smt/src/test/javaβ , executed with Maven and JUnit 5.
The scenarios under tests/casesβ mirror the Kafnus functional tests (same paths, expected_*.json, setup.sql and description.txt), but their input.json holds the exact messages Kafnus NGSI publishes to the processed topics for that scenario ("type": "raw"), recorded from a Kafnus run. This way they check the contract between both components without running Orion or Kafnus NGSI. When the Kafnus NGSI output changes, record the scenarios again and copy them over the existing ones.
Scenarios removed from Kafnus must be deleted here by hand. Scenarios specific to Kafnus Connect can still be written by hand with the postgis and mongo message types.
HTTP scenarios need the Kafnus Connect container to reach the HTTP mock started by pytest on the host (172.17.0.1:3333), so a host firewall must allow it.
Coverage is generated with JaCoCo for this SMT module and published to Coveralls from CI. It merges the SMT unit tests (Java) with the Python E2E suite: when KAFNUS_TESTS_COVERAGE=true, the E2E tests attach the JaCoCo agent to the Kafnus Connect container and write the execution data to tests/coverage/jacoco-e2e.exec when the stack stops.
To get the merged report locally (Maven must run with JDK 17, as the image does: JaCoCo silently ignores execution data when class files differ, and other JDKs compile different bytes):
cd tests && KAFNUS_TESTS_COVERAGE=true pytest -s test_pipeline.py && cd ..
mvn -f src/kafnus-connect-smt/pom.xml verify jacoco:merge@merge-e2e jacoco:report@report-merged
# Report in src/kafnus-connect-smt/target/site/jacoco-merged/index.html
src/kafnus-connect-smt/ covering JDBC and MongoDB sinks:
HeaderRouter: dynamic SQL routing for JDBC (schema/table resolution)MongoNamespacePrefix: MongoDB database/collection prefixing/usr/share/java/.For deeper technical details about how Kafnus Connect is configured, built, and extended β including:
π See Technical Configuration Guideβ
There is a way to change log level of kafnus-connect using logger API.
For example to change from default (info) to debug some clases related with http connector like HttpSinkTask, AbstractHttpSender, BasicAuthHttpSender you can use:
curl -s -X PUT -H "Content-Type: application/json" \
http://localhost:8083/admin/loggers/io.aiven.kafka.connect.http.HttpSinkTask \
-d '{"level":"DEBUG"}' | jq
curl -s -X PUT -H "Content-Type: application/json" \
http://localhost:8083/admin/loggers/io.aiven.kafka.connect.http.sender.AbstractHttpSender \
-d '{"level":"DEBUG"}' | jq
curl -s -X PUT -H "Content-Type: application/json" \
http://localhost:8083/admin/loggers/io.aiven.kafka.connect.http.sender.BasicAuthHttpSender \
-d '{"level":"DEBUG"}' | jq
To change from default (info) to debug clases related with jdbc connector:
curl -s -X PUT -H "Content-Type: application/json" \
http://localhost:8083/admin/loggers/io.confluent.connect.jdbc.sink.JdbcSinkTask \
-d '{"level":"DEBUG"}' | jq
π§ Project structure note
This repository is part of the Kafnus ecosystemβ :
The list of contributors to the Kafnus-Connect project can be found in
CONTRIBUTORS.mdβ .
Content type
Image
Digest
sha256:36f879774β¦
Size
1.1 GB
Last updated
about 17 hours ago
docker pull telefonicaiot/kafnus-connect