Skip to main content
Version: V2-Next

Apache NiFi integration

Apache NiFi is the pipeline engine of the platform. ADR 047 records that decision.

This page shows how the platform uses NiFi. It describes the installed system, and it describes how a Pipeline becomes a NiFi flow. The last section lists the requirements that are still open.

Apache NiFi in the platform

The important principle​

Users do not operate NiFi. The platform generates the flow.

A user models a Pipeline in the portal. The Pipeline is a graph of nodes. The config-adapter-nifi component reads that graph and writes a NiFi flow snapshot. It then sends the snapshot to NiFi through the REST API.

NiFi is an internal runtime. It is not a user interface of the platform.

Installation​

One NiFi node runs in the namespace. There is no cluster. High availability is disabled on purpose, because a second node needs a per-pod host name that the current chart does not give.

PropertyValue
Nodes1
StateKubernetes leases. ZooKeeper is disabled.
AuthenticationOIDC against Keycloak, client nifi
IngressDisabled. The service is ClusterIP.
Custom componentsNAR files from /opt/civitas/nars
MetricsAn exporter publishes metrics. A ServiceMonitor collects them.
Resources, development1 GiB to 2 GiB memory, 0.5 to 1 CPU
Resources, production8 GiB to 12 GiB memory, 2 to 4 CPU, JVM heap 4 g to 6 g

A Helm hook job runs after the installation and after an upgrade. The job creates two identities in NiFi and gives them their access policies: the human administrator, and the service account of config-adapter-nifi.

An init container builds the MQTT truststore before NiFi starts.

How a Pipeline becomes a NiFi flow​

The node vocabulary​

The adapter accepts six node kinds. Each kind has one role.

KindRole
sourceEmits the data. Exactly one per Pipeline.
mappingTransforms the data. Any number per Pipeline.
sinkWrites the data. Exactly one per Pipeline.
cronSchedules the source. It is not part of the data path.
start, endStructural anchors. They carry no data.

The CORE Pipeline schema also has filter, enrich and split. The adapter does not know them. If a graph puts one of them into the data path, the deployment fails.

The stage chain​

The adapter builds this chain:

source → [convert] → [mapping …] → sink

The convert stage is structural. A user cannot model it. The adapter inserts it only when the source emits a raw payload and the sink needs records.

StageNiFi components
MQTT sourceConsumeMQTT subscribes to one topic filter. It delivers the raw SensorThings envelope.
SQL sourceQueryDatabaseTableRecord reads a table over a connection pool. It emits records, and it supports a cron.
convertConvertRecord turns a raw payload into records.
mappingA chain of UpdateRecord processors. NiFi allows one replacement strategy per processor, so a Mapping that mixes constants with record paths needs more than one processor.
PostGIS sinkPutDatabaseRecord over the connection pool of the platform. A geometry value arrives as WKT, and the server parses it.
FROST sinkPutFrostRecord. This is a custom processor in the Civitas FROST NAR. It writes each record as one atomicity group of a JSON batch request.

One process group per Pipeline​

The adapter creates one process group for each Pipeline. The name of the group is pipeline- and the UUID of the Pipeline. The group is a direct child of the root process group.

The adapter finds its own groups by this name. A group with a different name is not managed, and the adapter does not touch it.

The deployment sequence​

  1. Get an OIDC token.
  2. Read the identifier of the root process group.
  3. Upload the flow snapshot.
  4. Push the sensitive properties into the controller services.
  5. Enable the controller services.
  6. Start the process group.

A group with the same name is stopped and deleted first, so a second deployment gives the same result as the first. A deletion of a group that does not exist is also a success.

A network error and an HTTP 5xx are retryable. An HTTP 4xx is fatal.

Only approved components​

The adapter mints every component from a curated list of fragments. It never generates a scripting processor. A Pipeline graph cannot add a processor that is not on the list.

The list of stages is wired by hand when the adapter starts. It is not discovered from the classpath, because an open registration would let any JAR add a component without a review.

This answers the third attacker in the threat model. An administrative user writes a Pipeline graph, not NiFi components. The user cannot put code into the runtime.

Secrets​

The flow snapshot holds no secret. The adapter collects the sensitive properties, and it sends them to the controller services after the upload. A secret is therefore never part of a stored flow definition.

Status and monitoring​

A monitor in the adapter polls NiFi at a fixed interval. It reads the processor status of each managed process group, and it reads the NiFi bulletin board.

The monitor publishes the state of the Pipeline back to the platform. A Pipeline that failed must be healthy in three consecutive rounds before the monitor reports a recovery.

Isolation​

A process group is a logical boundary. It is not a technical boundary. All process groups run in one JVM, and they share memory and thread pools.

The platform accepts this, for two reasons:

  • A user cannot put code into a process group. The curated component list prevents it.
  • NiFi access policies control the process groups. The bootstrap job gives the adapter identity and the administrator identity their policies.