Architekturentwurf: XML-Verarbeitungspipeline

grober Datenfluss von Eingangsdatei bis Ergebnisdatei
Kapitel 1 · Pipeline & Prozessschritte
Datenfluss von Eingangsdatei bis Ergebnisdatei, inkl. Details je Schritt
Prozess 1 · Empfang & Zerlegung
XML-Datei
Eingang
Ingestion
Empfang, ID vergeben, Status anlegen
Parser / Zerlegung
Streaming-Parser, fachliche Einheiten
parallele Weiterverarbeitung je Einheit (Producer-Consumer)
Consumer 1
Validierung / Transformation
Consumer 2
Validierung / Transformation
Consumer N
Validierung / Transformation
DB
Rohdaten / zerlegte Einheiten
Prozess 2 · Verarbeitungslogik
Queue
Entkopplung, Retry
Verarbeitungslogik / Worker
fachliche Logik, skalierbar
DB
Ergebnisse
Prozess 3 · Ausgabe-Generierung
XML-Generator
Ergebnisdaten → XML (Templating/XSLT)
XML-Datei
Ergebnis / Ausgang
Status- & Fehlertracking begleitet alle Stufen — pro Datei & pro Einheit, inkl. Dead-Letter/Retry
Datei
Datenbank
Message-Queue
Service/Komponente

Beschreibung der Schritte

1
Ingestion (Eingang)
Die große XML-Datei kommt über einen definierten Eingangskanal an, z.B. SFTP-Ordner, Objektspeicher (S3/Blob) mit Event-Trigger oder eine Message-Queue. Der Ingestion-Service nimmt die Datei entgegen, vergibt eine eindeutige Verarbeitungs-ID (Correlation-ID) und legt einen Eintrag in einer Statustabelle an ("empfangen"). Ab hier ist die Datei im System nachverfolgbar.

Besonderheit bei FTP/SFTP als Quelle, wenn Dateien nach der Verarbeitung nicht sofort gelöscht werden: Damit nicht bei jedem Poll erneut alle Dateien heruntergeladen und verarbeitet werden, sollte eine Datei entweder auf dem Server in einen Ordner wie "processed/" verschoben bzw. umbenannt werden (sofern Schreibrechte bestehen), oder – falls die Dateien unverändert liegen bleiben müssen – ein Abgleich der Verzeichnis-Metadaten (Dateiname, Größe, Änderungsdatum) gegen die bestehende Statustabelle erfolgen, bevor eine Datei heruntergeladen wird. Das Verschieben/Umbenennen erfolgt dabei praktischerweise bereits am Ende von Prozess 1, sobald Parsing und Persistierung in die DB erfolgreich abgeschlossen sind – nicht erst nach Verarbeitungslogik und XML-Generierung, da diese als eigene, entkoppelte Prozessketten ohnehin nur noch aus der DB lesen und die Quelldatei ab diesem Zeitpunkt nicht mehr benötigen. Das Verzeichnis-Listing selbst ist günstig, teuer ist nur das erneute Herunterladen und Verarbeiten. Ergänzend hilft ein Wasserzeichen (letztes erfolgreich verarbeitetes Änderungsdatum), um die Vergleichsmenge klein zu halten, ein Claim-Mechanismus bei mehreren parallelen Ingestion-Workern, sowie eine idempotente Persistenz als Rückversicherung gegen versehentliche Doppelverarbeitung.
2
Parser / Zerlegung
Weil die Dateien groß sind, wird ein Streaming-Parser eingesetzt (z.B. SAX/StAX in Java oder iterparse in Python), statt die Datei komplett in den Speicher zu laden. Der Parser zerlegt die XML-Struktur in fachlich sinnvolle Einheiten (z.B. einzelne Datensätze) und reicht diese einzeln zur Weiterverarbeitung weiter, statt auf das Ende der Datei zu warten.

Das Einlesen selbst bleibt dabei sequenziell, da der Parser Byte für Byte durch den Tokenstream läuft, um Element-Grenzen zu erkennen. Sobald der Parser jedoch eine fachliche Einheit vollständig erkannt hat (z.B. ein komplettes Element mit allen Kindelementen), kann die Weiterverarbeitung dieser Einheit – Validierung, Transformation, Schreiben in die DB – parallelisiert werden. Das ist ein klassisches Producer-Consumer-Muster: der Parser-Thread "produziert" laufend fertige Einheiten, während mehrere Consumer sie parallel abarbeiten.
3
DB – Rohdaten
Die zerlegten Einheiten werden persistiert, idealerweise per Batch-Insert statt Einzel-Inserts, um bei großen Mengen performant zu bleiben. Fachliche Felder liegen relational vor; falls Teile der ursprünglichen XML-Struktur erhalten bleiben müssen, kann ergänzend ein Blob-/Dokumentenspeicher für Roh-Fragmente sinnvoll sein.
4
Queue (Entkopplung)
Zwischen Persistenz und Verarbeitungslogik liegt eine Message-Queue (z.B. Kafka, RabbitMQ). Sie entkoppelt die Schritte zeitlich, ermöglicht horizontale Skalierung der nachfolgenden Worker und macht die Pipeline robuster gegenüber Lastspitzen – inklusive Retry-Mechanismus und Dead-Letter-Queue für dauerhaft fehlschlagende Einheiten.

Offene Frage dabei: Wer erzeugt die Nachricht und legt sie in die Queue? Direktes Dual-Write (erst DB-Insert, dann separat Queue-Publish im selben Prozessschritt) ist riskant, da beides unabhängige Systeme ohne gemeinsame Transaktion sind – schlägt einer der beiden Schritte fehl, sind DB und Queue inkonsistent.

Empfohlen: Transactional-Outbox-Pattern. Der schreibende Consumer legt in derselben DB-Transaktion, in der er die fachlichen Daten schreibt, zusätzlich einen Eintrag in einer Outbox-Tabelle an – dadurch atomar, entweder beides oder nichts. Ein separater, schlanker Relay-Prozess liest laufend diese Outbox-Tabelle (Polling) und veröffentlicht die Einträge in die Queue, markiert sie danach als versendet.

Alternative ohne echten Message-Broker: Ein Scheduler fragt periodisch die DB nach Datensätzen mit Status "roh, noch nicht verarbeitet" ab (Pull statt Push). Technisch keine Queue im klassischen Sinn, aber operativ am einfachsten und für moderaten Durchsatz oft ausreichend.

Anmerkung: Setzt man hierfür z.B. Spring Batch ein, kann die Parallelisierung intern über Multi-Threaded Step/Partitioning erfolgen – dann lässt sich die Queue tatsächlich einsparen. Vorteil: deutlich weniger Infrastruktur, kein Broker-Betrieb, eingebaute Restart-/Retry-Logik über das Job-Repository. Nachteil: Die Parallelisierung skaliert nur innerhalb einer JVM/eines Rechners – für echte horizontale Skalierung über mehrere Instanzen braucht es eigene Koordination (z.B. "SELECT ... FOR UPDATE SKIP LOCKED"), was eine Queue mit Consumer-Groups von Haus aus mitbringt. Zudem arbeitet ein Scheduler in Zyklen (Batch-Latenz), während eine Queue sofort reagiert – relevant, falls nahezu Echtzeit-Verarbeitung oder mehrere unabhängige Konsumenten der Ereignisse gebraucht werden. Für eine Instanz mit moderatem Durchsatz und tolerierbarer Batch-Latenz ist die Spring-Batch-Variante eine sinnvolle Vereinfachung.
5
Verarbeitungslogik / Worker
Ein oder mehrere Worker-Services lesen die persistierten Daten und führen die eigentliche fachliche Verarbeitungslogik durch. Durch die Entkopplung über die Queue lassen sich mehrere Worker-Instanzen parallel betreiben, um den Durchsatz bei Bedarf zu erhöhen.

Anmerkung: Als Worker eignet sich z.B. eine Message-Driven Bean (MDB), falls ein Java-EE/Jakarta-EE-Container (WildFly, Open Liberty o.ä.) im Einsatz ist – der Container ruft sie automatisch bei eingehender Nachricht auf und kann die Verarbeitung transaktional mit dem DB-Write verknüpfen. Bei Spring/Spring Boot ist das funktionale Gegenstück eine mit `@JmsListener` bzw. `@KafkaListener` annotierte Methode; Skalierung erfolgt in beiden Fällen über die Anzahl gleichzeitiger Consumer-Instanzen, nicht über zusätzliche Prozesse.
6
DB – Ergebnisse
Die berechneten Ergebnisse werden in eigenen Ergebnistabellen gespeichert, verknüpft mit den ursprünglichen Datensätzen über die Correlation-ID. So bleibt nachvollziehbar, welches Ergebnis zu welcher Eingangsdatei bzw. welchem Datensatz gehört.
7
XML-Generator
Ein Service liest die Ergebnisdaten aus der Datenbank und baut daraus die neue XML-Struktur, z.B. per Templating/XSLT oder Objekt-zu-XML-Mapping (etwa JAXB). Auch hier empfiehlt sich ein streamender Ansatz, wenn die Ergebnisdateien wieder groß werden können.

Mit "Templating" ist gemeint, dass die Ausgabestruktur als Vorlage mit Platzhaltern vordefiniert wird und eine Template-Engine (z.B. FreeMarker, Velocity) diese zur Laufzeit mit den Ergebnisdaten befüllt – im Unterschied zum programmatischen Aufbau der XML im Code. Beispiel einer solchen Vorlage:
<ergebnis> <id>${id}</id> <wert>${berechneterWert}</wert> <#list positionen as pos> <position>${pos.name}</position> </#list> </ergebnis>
Zum Vergleich: XSLT transformiert eine bestehende XML-Struktur über deklarative Regeln in eine andere (sinnvoll, wenn die Ergebnisdaten selbst schon als XML vorliegen), während JAXB (Objekt-zu-XML-Mapping) ganz ohne Vorlage auskommt – annotierte Java-Klassen werden direkt in XML serialisiert.
8
Ausgang
Die erzeugte XML-Ergebnisdatei wird an einem definierten Ausgangspunkt abgelegt (Ordner, Objektspeicher, SFTP-Ausgang) und der Verarbeitungsstatus wird auf "abgeschlossen" gesetzt. Optional werden nachgelagerte Systeme über Fertigstellung benachrichtigt.
Querschnittlich: Status- & Fehlertracking
Über alle Schritte hinweg begleitet eine Statustabelle jede Datei und – wo sinnvoll – jede einzelne Einheit darin. So lässt sich nachvollziehen, in welchem Schritt sich eine Datei befindet, wo Fehler aufgetreten sind und ob eine Verarbeitung wiederholt werden kann, ohne Duplikate zu erzeugen (Idempotenz). Fehlerhafte Einheiten landen in einer Dead-Letter-Queue statt die gesamte Datei zu blockieren.
Kapitel 2 · Architekturmuster
Empfehlungen für den internen Aufbau der einzelnen Prozessschritte
Verarbeitung
1
Pipeline (Chain-of-Responsibility) + Strategy für Varianten
Statt dass eine Service-Klasse mehrere Verarbeitungsklassen einzeln und explizit nacheinander aufruft, empfiehlt sich das Pipeline- bzw. Chain-of-Responsibility-Muster: Ein gemeinsames Interface für einen Verarbeitungsschritt, die konkreten Schritte (z.B. Validierung, Anreicherung, Berechnung, Ergebnis-Mapping) werden in einer Liste aufgereiht, durch die eine Pipeline nacheinander durchläuft.
interface VerarbeitungsSchritt<T> { T verarbeite(T einheit); } class ValidierungsSchritt implements VerarbeitungsSchritt<FachlicheEinheit> { ... } class AnreicherungsSchritt implements VerarbeitungsSchritt<FachlicheEinheit> { ... } class BerechnungsSchritt implements VerarbeitungsSchritt<FachlicheEinheit> { ... } class VerarbeitungsPipeline { private final List<VerarbeitungsSchritt<FachlicheEinheit>> schritte; FachlicheEinheit verarbeite(FachlicheEinheit einheit) { for (var schritt : schritte) { einheit = schritt.verarbeite(einheit); } return einheit; } } @Service class VerarbeitungsService { private final VerarbeitungsPipeline pipeline; private final ErgebnisRepository repo; @Transactional void verarbeiteEinheit(Long id) { var einheit = repo.findById(id); var ergebnis = pipeline.verarbeite(einheit); repo.speichereErgebnis(ergebnis); } }
Vorteil: jeder Schritt ist einzeln testbar, neue Schritte lassen sich hinzufügen ohne bestehende Klassen anzufassen (Open/Closed-Prinzip), und die Reihenfolge ist an einer Stelle sichtbar statt über Methodenaufrufe verstreut.

Falls die eigentliche Berechnung je nach Art der fachlichen Einheit unterschiedlich abläuft (z.B. verschiedene Dokument- oder Datensatztypen mit jeweils eigener Berechnungsvorschrift): zusätzlich das Strategy-Pattern statt if/else- oder switch-Kaskaden. Ein Interface, z.B. BerechnungsStrategy, mit mehreren Implementierungen je Typ, ausgewählt über eine Registry/Factory – in Spring z.B. über eine injizierte List<BerechnungsStrategy>, aus der die passende Implementierung über eine unterstuetzt(typ)-Methode ermittelt wird, ganz ohne zentrale Verzweigungslogik, die bei jedem neuen Typ wächst.

Anmerkung: Falls die Berechnung fachlich sehr komplex und reichhaltig mit Geschäftsregeln ist, lohnt sich zusätzlich ein reicheres Domänenmodell (DDD-Gedanke) – die Regeln direkt auf den fachlichen Objekten abbilden statt alles rein prozedural in Verarbeitungsklassen zu packen, damit die Objekte nicht zu reinen Datencontainern verkommen ("Anemic Domain Model"). Für einen groben Entwurf reicht Pipeline + Strategy aber völlig aus.
2
Explizites Statusmodell statt reiner In-Memory-Pipeline
Alternative bzw. Ergänzung zur reinen In-Memory-Pipeline: Jeder Satz bekommt einen expliziten, persistierten Status (Enum-Spalte) mit fest definierten, erlaubten Übergängen statt eines beliebigen Freitext-Felds – z.B. EMPFANGEN → PERSISTIERT → VALIDIERT bzw. AUSGESTEUERT → BERECHNET → XML_ERZEUGT → ABGESCHLOSSEN. Die Validierung ist dann der erste Übergang: im positiven Fall nach VALIDIERT, im negativen nach AUSGESTEUERT – nur im Zustand VALIDIERT wird der Satz für die nächste Stufe sichtbar. Technisch reicht dafür meist eine Status-Spalte plus eine Historie-/Event-Tabelle (alter Status, neuer Status, Zeitstempel, ggf. Fehlermeldung); bei komplexerer Zustandslogik (Wiederholungen, parallele Teilzustände, manuelle Eingriffe) lohnt sich eine echte State-Machine-Bibliothek wie Spring State Machine. Vorteil: Für eine GUI ergibt sich daraus fast von selbst ein Bild – eine Übersicht, wie viele Sätze aktuell in welchem Zustand stehen, und pro Satz die komplette Übergangs-Historie zum Nachvollziehen. Das deckt sich mit dem bereits geplanten Status-/Fehlertracking, formalisiert es aber zu einem echten Zustandsmodell. Zu beachten: Je mehr Zustände, desto wichtiger ein vorab durchdachtes Zustandsdiagramm mit klar erlaubten Übergängen, sonst driftet das mit der Zeit in unübersichtlichen "Status-Wildwuchs" ab.

Mehr Zustände bedeuten dabei nicht automatisch mehr Infrastruktur-Komponenten (MDBs/Listener). Zustand ist ein günstiges Datenmodell-Konzept, eine MDB/Queue dagegen eine teurere Infrastruktur-Komponente, die deployt, überwacht und skaliert werden muss – beides muss nicht 1:1 mitwachsen. Die MDB sollte an den bewusst gewählten Entkopplungsgrenzen sitzen (den bestehenden Prozessgrenzen 1 → 2 → 3), nicht an jedem einzelnen Zwischenzustand; innerhalb eines Consumer-Aufrufs kann ein Satz problemlos durch mehrere Zwischenzustände wandern, ohne weiteren Nachrichten-Hop. Für die Auswahl des passenden Handlers je Zustand bietet sich – analog zur Strategy oben – eine Registry an: ein Interface ZustandsHandler mit einer Implementierung pro Zustand/Übergang, ausgewählt über eine Map von Zustand auf Handler. Ein neuer Zustand bedeutet dann nur eine neue Handler-Klasse plus Registry-Eintrag, keine neue Queue, kein neuer Listener.

Ausnahme: Soll ein bestimmter Übergang bewusst unabhängig skalieren oder isoliert überwacht werden (z.B. eine rechenintensive Berechnung getrennt von der Validierung), ist eine eigene Queue/MDB dafür sinnvoll – das sollte dann aber eine bewusste Skalierungs-/Isolationsentscheidung sein, keine automatische Folge zusätzlicher Zustände.
Validierung
1
Fachliche Fehler getrennt von technischen Fehlern behandeln
Technische Fehler (DB nicht erreichbar, NullPointerException, Timeout) sind unerwartet – hier ist eine Exception mit Transaktions-Rollback und Retry/Dead-Letter-Queue richtig. Fachliche Fehler (Daten in der XML nicht stimmig, Pflichtfeld fehlt, Referenz auf unbekannten Stammdatensatz) sind dagegen kein Programmfehler, sondern ein erwartetes, gültiges Ergebnis der Validierung und sollten nicht als Exception mit Rollback modelliert werden. Sonst gehen zwei Dinge schief: Der Rollback wirft auch die Information weg, dass und warum der Satz abgelehnt wurde, und bei Retry/Redelivery über die Queue droht das klassische "Poison Message"-Problem – ein Satz, der nie gültig wird, aber trotzdem endlos erneut zugestellt wird.

Sauberer Ansatz: Die Validierung liefert kein Exception, sondern ein Ergebnis-Objekt zurück (Notification Pattern) – z.B. eine Liste von Validierungsfehlern mit Feld, Regel und Meldungstext. Das passt direkt zum Pipeline-Muster oben: Der ValidierungsSchritt ist der erste Schritt in der Pipeline und gibt entweder "OK" oder "fachlich fehlerhaft plus Fehlerliste" zurück. Bei "fachlich fehlerhaft" überspringt der orchestrierende Service die nachfolgenden Schritte (Verarbeitungslogik, Berechnung), committet den Satz aber ganz normal – nur mit eigenem Status wie "ausgesteuert"/"fachlich abgelehnt" statt "erfolgreich verarbeitet", zusammen mit den strukturierten Fehlerdetails. Wichtig fürs Statusmodell: technischer Fehler = retry-fähig, fachlicher Fehler = terminaler, aber gültiger Zustand, der nicht wie eine Störung eskaliert werden sollte.

Einreicher informieren: Dafür bietet sich derselbe Mechanismus an, der ohnehin für die Ergebnis-XML existiert (Prozess 3, XML-Generator). Ein zusätzlicher bzw. alternativer Rückmeldungs-/Ablehnungs-XML-Typ listet pro Satz auf, welche Regel verletzt wurde und warum – der XML-Generator liest dafür die Sätze mit Status "ausgesteuert" statt "erfolgreich berechnet" und erzeugt daraus eine Ablehnungsquittung für den Einreicher. Strukturell derselbe Prozess wie die normale Ergebnisdatei, nur mit anderem Template und anderer Datenquelle. Bei Bedarf ergänzend eine E-Mail-Benachrichtigung oder ein Reporting/Dashboard für die Fachabteilung, z.B. bei gehäuften Ablehnungen.
Marken- und Produktnamen
In diesem Dokument genannte Produkt-, Firmen- und Markennamen — u. a. Apache Kafka, RabbitMQ, Amazon S3, Java, Jakarta EE, Python, JAXB, Spring, Spring Boot, Spring Batch, Spring State Machine, WildFly, Open Liberty, Apache FreeMarker, Apache Velocity — sind Marken bzw. eingetragene Marken der jeweiligen Rechteinhaber. Sie werden hier ausschließlich zu illustrativen und erklärenden Zwecken verwendet; eine Zugehörigkeit, Empfehlung oder Zusammenarbeit mit den jeweiligen Unternehmen ist damit nicht verbunden.