Skip to main content
Version: 2.0.0

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.

Apache NiFi in the platform

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.

PropertyValue
Nodes1
StateKubernetes leases. ZooKeeper is disabled.
AuthenticationOIDC 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
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 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.

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.

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