Apache NiFi integration
Apache NiFi is the pipeline engine of the platform. ADR 047 records that decision.
This page shows how the platform uses NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows.. It describes the installed system, and it describes how a Pipeline becomes a NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. flow. The last section lists the requirements that are still open.
The important principle
Users do not operate NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows.. 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 NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. flow snapshot. It then
sends the snapshot to NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. through the REST API.
NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. is an internal runtime. It is not a user interface of the platform.
Installation
One NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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 KeycloakKeycloakAn open-source Identity and Access Management (IAM) solution providing SSO and OAuth2/OpenID Connect flows. In CIVITAS/CORE it is used to authenticate users and issue JWTs., 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 NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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 NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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 | NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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. NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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 NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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 NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. at a fixed interval. It reads the processor status of each managed process group, and it reads the NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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.
- NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. access policies control the process groups. The bootstrap job gives the adapter identity and the administrator identity their policies.