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
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
| Anforderung | Umsetzung |
|---|---|
| Cron/Scheduler | CRON-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-Verarbeitung | ExecuteSQL + QueryDatabaseTable für Bulk-Reads, MergeContent für Batching. |
| Stream-Verarbeitung | Kontinuierlicher FlowFile-Fluss: Prozessoren verarbeiten FlowFiles sofort nach Eingang. |
Konnektoren
| Anforderung | Umsetzung |
|---|---|
| MQTT | ConsumeMQTT, PublishMQTT (v3.1.1 + v5) |
| HTTP | InvokeHTTP (ausgehend), ListenHTTP, HandleHttpRequest/Response (eingehend) |
| SQL | ExecuteSQL, QueryDatabaseTable, PutDatabaseRecord, PutSQL (JDBC: PostgreSQL, MySQL, MSSQL, etc.) |
| Dateien | GetFile, PutFile, FetchFile, diverse Record-Reader/Writer (CSV, JSON, Avro, Parquet) |
| S3 | ListS3, FetchS3Object, PutS3Object, DeleteS3Object (AWS SDK, kompatibel mit MinIO und anderen S3-kompatiblen Stores) |
| SFTP | ListSFTP, FetchSFTP, PutSFTP, GetSFTP (SSH-basierter Dateitransfer) |
Prozessoren
| Anforderung | Umsetzung |
|---|---|
| JSON-Transformation | JoltTransformJSON (deklaratives Mapping), EvaluateJsonPath, UpdateAttribute, ReplaceText |
| Mapping/Validierung | Record-basiertes Processing: ConvertRecord, ValidateRecord mit JSON Schema |
| Kontrollfluss | RouteOnAttribute, RouteOnContent (if/else), DistributeLoad (Fan-out), Funnels (Fan-in) |
| Persistierung | PutDatabaseRecord, PutSQL, InvokeHTTP, PublishMQTT |
| Fehlerbehandlung | Failure-Relationships auf Connections, Auto-Retry mit Penalty-Duration, Bulletin-Board für Alerts |
Performance & Skalierung
| Anforderung | Umsetzung |
|---|---|
| Horizontale Skalierung | 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: mehrere Nodes als StatefulSet. Automatische Lastverteilung. |
| Vertikale Skalierung | JVM Heap und Thread-Pool-Konfiguration pro Node. |
| Ressourceneffizienz | Mindestens 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
| Anforderung | Umsetzung |
|---|---|
| Pipeline-Isolation | Process Groups als logische Isolation. Alle Pipelines laufen in derselben JVM -- keine echte Prozess-Isolation. Memory-Isolation nur auf Applikationsebene. |
| Pipeline-to-Pipeline Isolation | 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. 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 Isolation | 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: Lese-/Schreib-/Execute-Rechte pro Process Group, pro User/Gruppe. |
| User-to-Platform Isolation | Policies verhindern Zugriff auf Root-Process-Group und System-Prozessoren. Einschränkung externer Verbindungen über Controller Services. |
| RBAC | Eingebaut: Feingranulare Policies auf Process Group, Processor, Connection, Controller Service-Ebene. User- und Group-Management via LDAP/OIDC. |
| Secrets | 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. 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
| Anforderung | Umsetzung |
|---|---|
| K8s ohne CRDs | Ja. 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. |
| Helm | Community 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. |
| Versionierung | Pipeline-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-gapped | Container-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:
| Komponente | Beschreibung | Aufwand |
|---|---|---|
| Configuration-Adapter | Adapter, 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-Mapping | Mapping 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-Integration | 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.-Metriken (JMX/Prometheus Reporter) in Civitas-Monitoring integrieren. | gering |
| UI-Integration | Entscheidung: 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
| Risiko | Bewertung |
|---|---|
| Hoher Ressourcenbedarf | 4-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-Isolation | Alle Pipelines teilen sich eine JVM. RBAC schützt auf Applikationsebene, aber kein Container-/Prozess-Grenze zwischen Pipelines. |
| UI-Abhängigkeit | 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 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ät | Steile 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ät | Kein offizielles Apache Helm Chart. Community Charts variieren in Qualität und Aktualität. |