Skip to main content
Version: 2.0-rc2

Architekturskizze: Apache NiFi

Dieses Dokument beschreibt, wie die Anforderungen mit Apache NiFi architektonisch umgesetzt werden.

Siehe auch: Vergleich mit Redpanda Connect | Architekturskizze Redpanda Connect

Architekturübersicht​

Apache NiFi Architektur

Kernprinzip: Ein 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.-Cluster (StatefulSet) mit mehreren Nodes.

Pipelines werden als Process Groups innerhalb des Clusters modelliert. Isolation erfolgt über 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.-interne Policies (RBAC), nicht über Container-Grenzen -- alle Process Groups teilen sich eine JVM.

NiFi-Konzepte​

Process Groups​

Eine Process Group ist NiFisApache 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. zentrale Organisationseinheit -- vergleichbar mit einem Namespace oder Ordner. Sie kapselt eine Menge von Prozessoren, Connections und ggf. verschachtelte Sub-Process-Groups zu einer logischen Einheit.

Für Civitas bedeutet das: Jedes DataSetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. bekommt eine eigene Process Group. Innerhalb dieser Process Group liegen alle Prozessoren, die zu diesem DataSetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. gehören (MQTT-Consumer, Transformationen, SQL-Writer, etc.). Bei Bedarf können weitere Sub-Process-Groups genutzt werden, um Pipelines innerhalb eines DatasetsDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. zu trennen.

Eigenschaften:

  • Zugriffssteuerung: Jede Process Group hat eigene Policies (View, Modify, Operate). Damit lässt sich steuern, welcher User/Gruppe welche Pipeline sehen, bearbeiten oder starten darf.
  • Versionierung: Process Groups können über eine externe Registry versioniert werden (Import/Export als Flow-Snapshot via 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. Registry oder JSON-Export).
  • Verschachtelung: Process Groups können beliebig geschachtelt werden, z.B. eine übergeordnete Gruppe pro DataPool mit Sub-Groups pro DataSetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions..
  • Input/Output Ports: Process Groups kommunizieren untereinander über explizite Input- und Output-Ports -- das macht Datenflüsse zwischen DataSetsDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. sichtbar und kontrollierbar.

Einschränkung: Process Groups sind eine logische Grenze, keine technische. Alle Process Groups laufen im selben JVM-Prozess, teilen sich Speicher und Thread-Pools. Es gibt keine Container- oder Prozess-Isolation zwischen ihnen.

FlowFiles​

Ein FlowFile ist die Dateneinheit, die durch 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. fließt. Es besteht aus:

  • Content: Die eigentlichen Nutzdaten (z.B. ein JSON-Dokument, eine CSV-Zeile, ein Binär-Blob).
  • Attributes: Key-Value-Metadaten (z.B. filename, mqtt.topic, http.status.code, oder eigene Attribute).

Jeder Prozessor empfängt FlowFiles, verarbeitet sie und gibt neue FlowFiles aus. 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. verwaltet den Content im Content Repository auf Disk, sodass auch große Payloads verarbeitet werden können, ohne den JVM-Heap zu belasten.

Controller Services​

Controller Services sind gemeinsam genutzte Ressourcen wie Datenbankverbindungspools (DBCPConnectionPool), SSL-Kontexte oder Record-Reader/Writer. Sie werden auf Process-Group-Ebene definiert und stehen allen Prozessoren innerhalb dieser Gruppe zur Verfügung. Für Civitas relevant: Ein JDBC-Connection-Pool pro DataSetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions.-Process-Group, der die Datenbankzugangsdaten kapselt.

Umsetzung der Anforderungen​

Pipeline-Varianten​

AnforderungUmsetzung
Cron/SchedulerCRON-driven Scheduling Strategy auf Prozessoren (z.B. GenerateFlowFile oder ExecuteSQL mit CRON-Timer).
Ereignisgetriebene Verarbeitung (MQTT)ConsumeMQTT-Prozessor subscribt auf Topics. Daten fließen als FlowFiles durch die Pipeline.
Ereignisgetriebene Verarbeitung (Webhook)ListenHTTP-Prozessor empfängt POST-Requests und erzeugt FlowFiles.
HTTP-Endpunkte (Request/Response)HandleHttpRequest + HandleHttpResponse-Prozessor-Paar für synchrones Request/Response.
Batch-VerarbeitungExecuteSQL + QueryDatabaseTable für Bulk-Reads, MergeContent für Batching.
Stream-VerarbeitungKontinuierlicher FlowFile-Fluss: Prozessoren verarbeiten FlowFiles sofort nach Eingang.

Konnektoren​

AnforderungUmsetzung
MQTTConsumeMQTT, PublishMQTT (v3.1.1 + v5)
HTTPInvokeHTTP (ausgehend), ListenHTTP, HandleHttpRequest/Response (eingehend)
SQLExecuteSQL, QueryDatabaseTable, PutDatabaseRecord, PutSQL (JDBC: PostgreSQL, MySQL, MSSQL, etc.)
DateienGetFile, PutFile, FetchFile, diverse Record-Reader/Writer (CSV, JSON, Avro, Parquet)
S3ListS3, FetchS3Object, PutS3Object, DeleteS3Object (AWS SDK, kompatibel mit MinIO und anderen S3-kompatiblen Stores)
SFTPListSFTP, FetchSFTP, PutSFTP, GetSFTP (SSH-basierter Dateitransfer)

Prozessoren​

AnforderungUmsetzung
JSON-TransformationJoltTransformJSON (deklaratives Mapping), EvaluateJsonPath, UpdateAttribute, ReplaceText
Mapping/ValidierungRecord-basiertes Processing: ConvertRecord, ValidateRecord mit JSON Schema
KontrollflussRouteOnAttribute, RouteOnContent (if/else), DistributeLoad (Fan-out), Funnels (Fan-in)
PersistierungPutDatabaseRecord, PutSQL, InvokeHTTP, PublishMQTT
FehlerbehandlungFailure-Relationships auf Connections, Auto-Retry mit Penalty-Duration, Bulletin-Board für Alerts

Performance & Skalierung​

AnforderungUmsetzung
Horizontale SkalierungNiFiApache 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.-Cluster: mehrere Nodes als StatefulSet. Automatische Lastverteilung.
Vertikale SkalierungJVM Heap und Thread-Pool-Konfiguration pro Node.
RessourceneffizienzMindestens 2-4 GB RAM pro 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 (JVM). Ein 2-Node-Cluster benötigt 4-8 GB nur für 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., auch ohne Pipelines.
Hochfrequente Daten (0.5-1s)Machbar, aber FlowFile-Overhead pro Nachricht. Für hohe Raten empfiehlt sich MergeContent (Micro-Batching).

Security & Isolation​

AnforderungUmsetzung
Pipeline-IsolationProcess Groups als logische Isolation. Alle Pipelines laufen in derselben JVM -- keine echte Prozess-Isolation. Memory-Isolation nur auf Applikationsebene.
Pipeline-to-Pipeline IsolationJedes DataSetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. wird als eigene Process Group mit eigenen Policies modelliert -- separate Security Domain. Pipelines innerhalb eines DataSetsDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. teilen sich eine Process Group und damit eine gemeinsame Domain. Datenflüsse zwischen DataSetsDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. sind nur über explizite Input/Output Ports möglich.
User-to-User IsolationNiFiApache 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.-interne Policies: Lese-/Schreib-/Execute-Rechte pro Process Group, pro User/Gruppe.
User-to-Platform IsolationPolicies verhindern Zugriff auf Root-Process-Group und System-Prozessoren. Einschränkung externer Verbindungen über Controller Services.
RBACEingebaut: Feingranulare Policies auf Process Group, Processor, Connection, Controller Service-Ebene. User- und Group-Management via LDAP/OIDC.
SecretsNiFiApache 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. Sensitive Properties (verschlüsselt in Flow-Config), Parameter Contexts mit externen Secret Providern (HashiCorp VaultHashiCorp VaultA secrets management system used as the platform's secrets backend. CIVITAS/CORE supports HashiCorp Vault (BSL) and its API-compatible open-source fork OpenBao. The platform does not implement its own encryption or secret storage., AWS Secrets Manager, Azure Key VaultHashiCorp VaultA secrets management system used as the platform's secrets backend. CIVITAS/CORE supports HashiCorp Vault (BSL) and its API-compatible open-source fork OpenBao. The platform does not implement its own encryption or secret storage., GCP Secret Manager). Prozessoren referenzieren Secrets via #{parameter_name} -- der Klartext ist nie in der Flow-Definition sichtbar.

Sicherheitshinweis: 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. bietet starke logische Isolation (RBAC), aber keine Prozess-/Container-Isolation zwischen Pipelines. Ein bösartiger Prozessor (z.B. ExecuteScript) kann potenziell auf Speicher anderer Pipelines zugreifen. Mitigation: Prozessoren mit Code-Ausführung (ExecuteScript, ExecuteGroovyScript, ExecuteStreamCommand) werden per 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.-Policy auf Root-Ebene auf Admin-Konten beschränkt. Verwaltende Nutzer erhalten nur Zugriff auf konfigurationsbasierte Prozessoren (JoltTransformJSON, RouteOnAttribute, etc.). Dies adressiert die Angreifer-Klasse "Authentifizierter Nutzer (verwaltend)" aus dem Threat Model.

Betrieb & Deployment​

AnforderungUmsetzung
K8s ohne CRDsJa. 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. als StatefulSet, kein Operator/CRD nötig.
HelmCommunity Helm Charts verfügbar (z.B. cetic/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.). Kein offizielles Apache Helm Chart.
VersionierungPipeline-Definitionen werden als Flow-Snapshots (JSON) in einer externen Registry versioniert. Der Configuration-Adapter importiert die aktuelle Version über die 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. REST API.
Air-gappedContainer-Images vorab ladbar. Keine Runtime-Abhängigkeiten. NARs (Plugins) sind im Image enthalten.

Eigenentwicklungsanteil​

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. bringt als Plattform mehr mit als Redpanda Connect, erfordert aber andere Anpassungen:

KomponenteBeschreibungAufwand
Configuration-AdapterAdapter, der Civitas-Pipeline-Definitionen 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. Process Groups übersetzt (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. REST API als Backend).mittel
RBAC-MappingMapping von Civitas-Rollen/DataSetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions.-Berechtigungen auf 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.-Policies und -Gruppen (ebenfalls im CA).gering
Monitoring-IntegrationNiFiApache 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.-Metriken (JMX/Prometheus Reporter) in Civitas-Monitoring integrieren.gering
UI-IntegrationEntscheidung: 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.-UI einbetten (iFrame/SSO) oder eigenes UI + 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.-REST-API?variabel

Risiken​

RisikoBewertung
Hoher Ressourcenbedarf4-8 GB RAM Baseline für einen 2-Node-Cluster (ohne Pipelines). Steigt mit Anzahl der Prozessoren. Widerspricht der Anforderung nach effizienter Ressourcennutzung bei ungenutzten Sandboxes.
Keine echte Pipeline-IsolationAlle Pipelines teilen sich eine JVM. RBAC schützt auf Applikationsebene, aber kein Container-/Prozess-Grenze zwischen Pipelines.
UI-AbhängigkeitNiFiApache 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. bietet eine starke UI. Code basierte Bereitstellung von Pipelines ist versioniert über Registry möglich. Aufwand für eigenen Editor ist zu prüfen, da eine Modellierungs-UI zumindest für Experten existiert
KomplexitätSteile Lernkurve. 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.-eigene Konzepte (FlowFiles, Provenance, Bulletin Board, Controller Services) erfordern Einarbeitung.
Helm-Chart-QualitätKein offizielles Apache Helm Chart. Community Charts variieren in Qualität und Aktualität.