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.
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.
| Property | Value |
|---|---|
| Nodes | 1 |
| State | Kubernetes leases. ZooKeeper is disabled. |
| Authentication | OIDC against Keycloak, client nifi |
| Ingress | Disabled. The service is ClusterIP. |
| Custom components | NAR files from /opt/civitas/nars |
| Metrics | An exporter publishes metrics. A ServiceMonitor collects them. |
| Resources, development | 1 GiB to 2 GiB memory, 0.5 to 1 CPU |
| Resources, production | 8 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.
| Kind | Role |
|---|---|
source | Emits the data. Exactly one per Pipeline. |
mapping | Transforms the data. Any number per Pipeline. |
sink | Writes the data. Exactly one per Pipeline. |
cron | Schedules the source. It is not part of the data path. |
start, end | Structural 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.
| Stage | NiFi components |
|---|---|
| MQTT source | ConsumeMQTT subscribes to one topic filter. It delivers the raw SensorThings envelope. |
| SQL source | QueryDatabaseTableRecord reads a table over a connection pool. It emits records, and it supports a cron. |
| convert | ConvertRecord turns a raw payload into records. |
| mapping | A 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 sink | PutDatabaseRecord over the connection pool of the platform. A geometry value arrives as WKT, and the server parses it. |
| FROST sink | PutFrostRecord. 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
- Get an OIDC token.
- Read the identifier of the root process group.
- Upload the flow snapshot.
- Push the sensitive properties into the controller services.
- Enable the controller services.
- 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.