Sign inSign up

telefonicaiot/kafnus-connect

By telefonicaiot

β€’Updated about 17 hours ago

Kafnus connect container

Image
0

7.9K

telefonicaiot/kafnus-connect repository overview

FIWARE Incubating Coverage Status

Coverage badge scope: custom Java SMTs (HeaderRouter and MongoNamespacePrefix), merging unit tests and Python E2E tests.

β πŸ›°οΈ Kafnus Connect

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⁠.


β βš™οΈ Overview

Kafnus Connect consumes processed NGSI events from Kafka topics (produced by Kafnus NGSI⁠) and persists them into target datastores or APIs.

⁠Supported sinks
  • πŸ—ΊοΈ PostGIS (via custom JDBC connector⁠)
    • Forked and extended to handle GeoJSON geometries and NGSI-specific data structures.
  • πŸ“¦ MongoDB
    • Official MongoDB Kafka connector for JSON document storage.
  • 🌐 HTTP

⁠🧱 Architecture

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.


β πŸš€ Usage

⁠Build locally
docker build -t telefonicaiot/kafnus-connect:latest .
⁠Run example
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.


⁠πŸ§ͺ Testing

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

⁠🧰 Configuration & Extensions

  • Custom SMTs are available in src/kafnus-connect-smt/ covering JDBC and MongoDB sinks:
    • HeaderRouter: dynamic SQL routing for JDBC (schema/table resolution)
    • MongoNamespacePrefix: MongoDB database/collection prefixing
  • New sinks can be added by extending the base image and adding plugins under /usr/share/java/.
  • Monitoring via Prometheus JMX Exporter is supported out of the box.

For deeper technical details about how Kafnus Connect is configured, built, and extended β€” including:

  • Environment setup and logging configuration
  • Plugin management and sink registration
  • Supported sinks and custom SMTs
  • Usage of EnvVarConfigProvider for connector configuration

πŸ‘‰ See Technical Configuration Guide⁠


β πŸ› οΈ Logging

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

β πŸ“š Documentation


🧭 Project structure note

This repository is part of the Kafnus ecosystem⁠:


⁠πŸ‘₯ Contributors

The list of contributors to the Kafnus-Connect project can be found in CONTRIBUTORS.md⁠.

Tag summary

Content type

Image

Digest

sha256:36f879774…

Size

1.1 GB

Last updated

about 17 hours ago

docker pull telefonicaiot/kafnus-connect