
Hallo, Habr!
Unser Unternehmen ist auf die Entwicklung von ERP-Softwarelösungen spezialisiert, wobei ein Großteil aus transaktionalen Systemen mit einer enormen Menge an Geschäft Logik und Dokumentenverkehr à la EAD besteht. Die modernen Versionen unserer Produkte basieren auf JavaEE-Technologien, aber wir experimentieren auch aktiv mit Mikrodiensten. Eines der problematischsten Bereiche solcher Lösungen ist die Integration verschiedener Teilsysteme, die zu benachbarten Domänen gehören. Integrationsaufgaben haben uns immer große Kopfschmerzen bereitet, unabhängig von den verwendeten Architekturstilen, technologischen Stacks und Frameworks. In letzter Zeit gibt es jedoch Fortschritte bei der Lösung solcher Aufgaben.
In dem Artikel, den ich Ihnen präsentiere, werde ich über die Erfahrungen und architektonischen Forschungen von NPO «Krista» in diesem Bereich berichten. Wir werden auch ein einfaches Beispiel für die Lösung einer Integrationsaufgabe aus der Sicht eines Anwendungsentwicklers betrachten und herausfinden, was hinter dieser Einfachheit steckt.
Haftungsausschluss
Die in dem Artikel beschriebenen architektonischen und technischen Lösungen basieren auf meinen persönlichen Erfahrungen im Kontext spezifischer Aufgaben. Diese Lösungen beanspruchen keine Universalität und können sich unter anderen Nutzungbedingungen als nicht optimal erweisen.
Was hat BPM damit zu tun?
Um diese Frage zu beantworten, müssen wir etwas tiefer in die spezifischen Anwendungsprobleme unserer Lösungen eintauchen. Der Großteil der Geschäft Logik in unserem typischen transaktionalen System besteht aus der Dateneingabe in die Datenbank über Benutzeroberflächen, der manuellen und automatisierten Prüfung dieser Daten, der Durchführung durch einen bestimmten Workflow, der Veröffentlichung in ein anderes System / Analyse-Datenbank / Archiv, der Erstellung von Berichten. Daher ist die Schlüssel-Funktion des Systems für die Kunden die Automatisierung ihrer internen Geschäftsprozesse.
Zur Vereinfachung verwenden wir in der Kommunikation den Begriff „Dokument“ als eine gewisse Abstraktion einer Datenmenge, die durch einen gemeinsamen Schlüssel verbunden ist, an den ein bestimmter Workflow „angehängt“ werden kann.
Wie sieht es aber mit der Integrationslogik aus? Schließlich entsteht die Integrationsaufgabe durch die Architektur des Systems, die nicht auf die Anforderungen des Kunden zugeschnitten ist, sondern durch ganz andere Faktoren beeinflusst wird:
- unter dem Einfluss des Conway-Gesetzes;
- durch die Wiederverwendung von Teilsystemen, die zuvor für andere Produkte entwickelt wurden;
- nach Entscheidung des Architekten, basierend auf den nicht-funktionalen Anforderungen.
Es gibt eine große Versuchung, die Integrationslogik von der Geschäftslogik des Haupt-Workflows zu trennen, um die Geschäftslogik nicht mit Integrationsartefakten zu belasten und dem Anwendungsentwickler die Notwendigkeit zu ersparen, sich mit den Besonderheiten der Systemarchitektur auseinanderzusetzen. Diese Herangehensweise hat zwar einige Vorteile, jedoch zeigt die Praxis ihre Ineffektivität:
- die Lösung von Integrationsproblemen tendiert oft zu den einfachsten Optionen in Form von synchronen Aufrufen aufgrund der begrenzten Erweiterungspunkte in der Implementierung des Haupt-Workflows (zu den Nachteilen synchroner Integration – weiter unten);
- Integrationsartefakte dringen dennoch in die Hauptgeschäftslogik ein, wenn Rückmeldungen aus einem anderen Teilsystem erforderlich sind;
- der Anwendungsentwickler ignoriert die Integration und kann sie leicht brechen, indem er den Workflow ändert;
- das System hört auf, aus der Sicht des Benutzers ein Ganzes zu sein, es werden 'Nähte' zwischen den Teilsystemen sichtbar, es entstehen überflüssige Benutzeroperationen, die die Datenübertragung von einem Teilsystem zum anderen einleiten.
Ein anderer Ansatz besteht darin, Integrationsinteraktionen als wesentlichen Bestandteil der Hauptgeschäftslogik und des Workflows zu betrachten. Damit die Anforderungen an die Qualifikation der Anwendungsentwickler nicht ins Unermessliche steigen, sollte die Erstellung neuer Integrationsinteraktionen einfach und unkompliziert erfolgen, mit minimalen Wahlmöglichkeiten für die Lösungsmethoden. Dies ist schwieriger zu realisieren, als es scheint: Das Werkzeug muss leistungsstark genug sein, um dem Benutzer eine ausreichende Anzahl an Anwendungsmöglichkeiten zu bieten, dabei jedoch verhindern, dass er sich "ins eigene Fleisch schneidet". Es gibt viele Fragen, die ein Ingenieur im Kontext von Integrationsaufgaben beantworten muss, über die sich der Anwendungsentwickler in seiner alltäglichen Arbeit jedoch keine Gedanken machen sollte: Transaktionsgrenzen, Konsistenz, Atomizität, Sicherheit, Skalierung, Last- und Ressourcenverteilung, Routing, Marshalling, Kontextverbreitung und -wechsel usw. Es sollten den Anwendungsentwicklern ausreichend einfache Lösungsvorlagen angeboten werden, die bereits Antworten auf all diese Fragen enthalten. Diese Vorlagen sollten ausreichend sicher sein: Die Geschäftslogik ändert sich sehr häufig, was die Fehleranfälligkeit erhöht, und die Kosten für Fehler müssen auf einem ausreichend niedrigen Niveau gehalten werden.
Aber was hat BPM damit zu tun? Es gibt doch viele Möglichkeiten, Workflows zu implementieren …
In unseren Lösungen ist tatsächlich eine andere Implementierung von Geschäftsprozessen sehr beliebt – über die deklarative Definition von Zustandsübergangsdiagrammen und die Anbindung von Handlern mit Geschäftslogik an die Übergänge. Dabei ist der Zustand, der die aktuelle Position des „Dokuments“ im Geschäftsprozess bestimmt, ein Attribut des „Dokuments“ selbst.

So sieht der Prozess zu Beginn des Projekts aus.
Die Beliebtheit einer solchen Umsetzung beruht auf der relativen Einfachheit und Geschwindigkeit bei der Erstellung linearer Geschäftsprozesse. Mit der fortwährenden Komplexität von Softwaresystemen wächst jedoch der automatisierte Teil des Geschäftsprozesses und wird komplizierter. Es entsteht die Notwendigkeit für Dekomposition, Wiederverwendung von Prozessbestandteilen sowie für das Verzweigen von Prozessen, damit jeder Zweig parallel ausgeführt werden kann. Unter diesen Bedingungen wird das Werkzeug unpraktisch, und das Zustandsübergangsdiagramm verliert an Informativität (Integrationsinteraktionen sind auf dem Diagramm überhaupt nicht abzubilden).

So sieht der Prozess nach mehreren Iterationen zur Klärung der Anforderungen aus
Eine Lösung für diese Situation war die Integration der Engine in einige Produkte mit den komplexesten Geschäftsprozessen. Kurzfristig hatte diese Lösung einen gewissen Erfolg: Es wurde möglich, komplexe Geschäftsprozesse umzusetzen und dabei ein ausreichend informatives und aktuelles Diagramm in der Notation .

Ein kleiner Teil eines komplexen Geschäftsprozesses
Auf lange Sicht erfüllte die Lösung jedoch nicht die Erwartungen: Der hohe Aufwand für die Erstellung von Geschäftsprozessen über visuelle Werkzeuge ließ keine akzeptablen Produktivitätskennzahlen erreichen, und das Werkzeug selbst wurde zu einem der unbeliebtesten unter den Entwicklern. Auch an der internen Struktur der Engine gab es Kritik, die zur Entstehung vieler 'Patches' und 'Kreuzstützen' führte.
Der Hauptvorteil der Anwendung von jBPM war das Bewusstsein für den Nutzen und die Nachteile des Vorhandenseins eines eigenen persistenten Zustands für ein Geschäftsprozess-Exemplar. Außerdem erkannten wir die Möglichkeit der Anwendung des prozessorientierten Ansatzes zur Umsetzung komplexer Integrationsprotokolle zwischen verschiedenen Anwendungen unter Verwendung asynchroner Interaktionen durch Signale und Nachrichten. Der vorhandene persistente Zustand spielt hier eine entscheidende Rolle.
Auf Grundlage des Gesagten lässt sich schließen: Der prozessorientierte Ansatz im Stil von BPM ermöglicht es uns, ein breites Spektrum an Aufgaben zur Automatisierung zunehmend komplexer werdender Geschäftsprozesse zu lösen, Integrationsaktivitäten harmonisch in diese Prozesse einzupassen und die Möglichkeit der visuellen Darstellung des umgesetzten Prozesses in der dafür geeigneten Notation zu bewahren.
Nachteile synchroner Aufrufe als Integrationsmuster
Unter synchroner Integration versteht man den einfachsten blockierenden Aufruf. Ein Teilsystem fungiert als Server-Seite und stellt eine API mit der benötigten Methode zur Verfügung. Das andere Teilsystem fungiert als Client-Seite und führt zum richtigen Zeitpunkt den Aufruf mit dem Warten auf das Ergebnis aus. Je nach Architektur des Systems können die Client- und Server-Seiten entweder in einer Anwendung und einem Prozess oder in verschiedenen platziert sein. Im zweiten Fall ist es erforderlich, eine bestimmte Implementierung von RPC anzuwenden und das Marshalling von Parametern und Ergebnis des Aufrufs zu gewährleisten.

Ein solches Integrationsmuster hat eine Vielzahl von Nachteilen, wird jedoch in der Praxis aufgrund seiner Einfachheit sehr häufig verwendet. Die Geschwindigkeit der Implementierung ist verlockend und führt dazu, dass es immer wieder unter "dringenden" Terminen eingesetzt wird, wobei die Lösung in technische Schulden übertragen wird. Oft kommt es auch vor, dass unerfahrene Entwickler es unbewusst anwenden, ohne sich der negativen Konsequenzen bewusst zu sein.
Neben der offensichtlichsten Erhöhung der Kopplung von Teilsystemen gibt es auch weniger offensichtliche Probleme mit der "Zerstreuung" und "Dehnung" von Transaktionen. Tatsächlich, wenn die Geschäftslogik Änderungen vornimmt, sind Transaktionen unvermeidbar, und Transaktionen blockieren ihrerseits bestimmte Ressourcen der Anwendung, die von diesen Änderungen betroffen sind. Das bedeutet, dass ein Teilsystem, bis es eine Antwort vom anderen erhält, die Transaktion nicht abschließen und die Blockierungen aufheben kann. Dies erhöht erheblich das Risiko verschiedener Effekte:
- Die Reaktionsfähigkeit des Systems geht verloren, Benutzer warten lange auf Antworten auf Anfragen;
- Der Server hört insgesamt auf, auf Benutzeranfragen zu reagieren, da der Pool der Threads überfüllt ist: Die meisten Threads sind aufgrund der Blockierung einer Ressource, die von der Transaktion benötigt wird, "blockiert";
- Es beginnen Deadlocks aufzutreten: Die Wahrscheinlichkeit ihres Auftretens hängt stark von der Dauer der Transaktionen, der Anzahl der in die Transaktion involvierten Geschäftslogik und der Blockierungen ab;
- Es treten Transaktionszeitüberschreitungsfehler auf;
- Der Server "stürzt ab" wegen OutOfMemory, wenn die Aufgabe die Verarbeitung und Änderung großer Datenmengen erfordert, und das Vorhandensein synchroner Integrationen erschwert die Aufteilung der Verarbeitung in leichtere Transaktionen.
Aus architektonischer Sicht führt die Verwendung blockierender Aufrufe bei der Integration zu einem Verlust der Kontrolle über die Qualität einzelner Teilsysteme: Es ist nicht möglich, die Zielwerte der Qualität eines Teilsystems unabhängig von den Qualitätswerten eines anderen Teilsystems zu gewährleisten. Wenn Teilsysteme von unterschiedlichen Teams entwickelt werden, stellt dies ein großes Problem dar.
Es wird noch interessanter, wenn die integrierten Teilsysteme sich in verschiedenen Anwendungen befinden und synchronisierte Änderungen von beiden Seiten vorgenommen werden müssen. Wie kann die Transaktionssicherheit dieser Änderungen gewährleistet werden?
Wenn Änderungen in getrennten Transaktionen vorgenommen werden, ist es erforderlich, eine zuverlässige Verarbeitung von Ausnahmen und Kompensationen sicherzustellen, was den Hauptvorteil synchroner Integrationen – die Einfachheit – vollständig negiert.
Es kommen auch verteilte Transaktionen in den Sinn, die wir jedoch in unseren Lösungen nicht verwenden: Es ist schwierig, die Zuverlässigkeit sicherzustellen.
„Saga“ als Lösung für das Transaktionsproblem
Mit der wachsenden Popularität von Mikrodiensten gewinnt .
Dieser Entwurf löst hervorragend die oben genannten Probleme langfristiger Transaktionen und erweitert die Möglichkeiten der Zustandsverwaltung des Systems aus der Sicht der Geschäftslogik: Eine Kompensation nach einer fehlgeschlagenen Transaktion kann das System nicht in den ursprünglichen Zustand zurückversetzen, sondern einen alternativen Datenverarbeitungsweg gewährleisten. Dadurch wird auch verhindert, dass erfolgreich abgeschlossene Schritte der Datenverarbeitung bei erneutem Versuch wiederholt werden.
Interessanterweise ist dieses Muster auch in monolithischen Systemen relevant, wenn es um die Integration schwach gekoppelter Teilsysteme geht und negative Effekte auftreten, die durch langfristige Transaktionen und entsprechende Ressourcensperrungen verursacht werden.
In Bezug auf unsere Geschäftsprozesse im BPM-Stil ist die Implementierung von „Sagas“ sehr einfach: Einzelne Schritte der „Saga“ können als Aktivitäten innerhalb des Geschäftsprozesses definiert werden, und der persistente Zustand des Geschäftsprozesses bestimmt auch den internen Zustand der „Saga“. Das heißt, wir benötigen keinen zusätzlichen Koordinationsmechanismus. Es wird nur ein Nachrichtenbroker mit Unterstützung für „at least once“-Garantien als Transport benötigt.
Doch auch diese Lösung hat ihren eigenen „Preis“:
- Die Geschäftslogik wird komplexer: Es müssen Kompensationen bearbeitet werden;
- Es wird erforderlich sein, auf volle Konsistenz zu verzichten, was für monolithische Systeme besonders empfindlich sein kann;
- Die Architektur wird etwas komplizierter, es entsteht ein zusätzlicher Bedarf an einem Nachrichtenbroker;
- Zusätzliche Mittel zur Überwachung und Verwaltung werden benötigt (obwohl das insgesamt sogar gut ist: Die Servicequalität des Systems wird steigen).
Für monolithische Systeme ist die Rechtfertigung der Verwendung von "Sagas" nicht so offensichtlich. Für Mikrodienste und andere SOA, bei denen wahrscheinlich bereits ein Broker vorhanden ist und volle Konsistenz bereits zu Beginn des Projekts geopfert wurde, kann der Nutzen der Anwendung dieses Musters die Nachteile erheblich überwiegen, besonders wenn eine praktische API auf der Ebene der Geschäftslogik vorhanden ist.
Kapselung der Geschäftslogik in Mikrodiensten
Als wir begannen, mit Mikrodiensten zu experimentieren, stellte sich die berechtigte Frage: Wo soll die domain-spezifische Geschäftslogik in Bezug auf den Service platziert werden, der die Persistenz der domain-spezifischen Daten gewährleistet?
Wenn man sich die Architektur verschiedener BPMS ansieht, erscheint es sinnvoll, die Geschäftslogik von der Persistenz zu trennen: eine Schicht von plattform- und domain-unabhängigen Mikrodiensten zu schaffen, die die Umgebung und den Container für die Ausführung der domain-spezifischen Geschäftslogik bilden, während die Persistenz der domain-spezifischen Daten in einer separaten Schicht aus sehr einfachen und leichtgewichtigen Mikrodiensten gestaltet wird. In diesem Fall orchestrieren die Geschäftsprozesse die Services der Persistenzschicht.

Dieser Ansatz hat einen großen Vorteil: Man kann die Funktionalität der Plattform beliebig erweitern, und die entsprechende Schicht der plattform-spezifischen Mikrodienste wird dadurch nur 'dicker'. Geschäftsprozesse aus jeder Domain erhalten sofort die Möglichkeit, die neue Funktionalität der Plattform zu nutzen, sobald sie aktualisiert wird.
Eine detaillierte Ausarbeitung hat erhebliche Nachteile dieses Ansatzes aufgezeigt:
- Der plattformbasierte Service, der die Geschäftslogik vieler Domains gleichzeitig ausführt, trägt hohe Risiken als einzelner Punkt des Fehlers. Häufige Änderungen der Geschäftslogik erhöhen das Risiko von Fehlern, die zu Störungen führen, die sich auf das gesamte System ausbreiten;
- Leistungsprobleme: Die Geschäftslogik arbeitet über eine enge und langsame Schnittstelle mit ihren Daten.
- Daten werden erneut marshalliert und durch den Netzwerk-Stack gepumpt;
- Der Domainservice gibt oft mehr Daten zurück, als für die Geschäftslogik zur Verarbeitung erforderlich ist, aufgrund unzureichender Parametrisierungsmöglichkeiten der Anfragen auf der Ebene der externen API des Dienstes;
- Mehrere unabhängige Teile der Geschäftslogik können dieselben Daten erneut abfragen (dies kann durch das Hinzufügen von Sitzungsbestandteilen, die Daten cachen, abgemildert werden, was jedoch die Architektur zusätzlich kompliziert und Probleme mit der Aktualität der Daten und der Invalidierung des Caches schafft);
- Transaktionsprobleme:
- Geschäftsprozesse mit persistentem Zustand, dessen Speicherung vom Plattformdienst verwaltet wird, werden von den Domain-Daten in Konflikt gebracht, und einfache Lösungen für dieses Problem sind nicht in Sicht;
- Verlagerung der Sperrung von Domain-Daten außerhalb der Transaktion: Wenn die Domain-Geschäftslogik Änderungen vornehmen muss, nachdem sie die Richtigkeit der aktuellen Daten überprüft hat, muss die Möglichkeit einer konkurrierenden Änderung der verarbeiteten Daten ausgeschlossen werden. Eine externe Datensperre kann helfen, das Problem zu lösen, bringt jedoch zusätzliche Risiken mit sich und verringert die Gesamtnutzerfreundlichkeit des Systems;
- Zusätzliche Komplikationen bei der Aktualisierung: In einigen Fällen müssen der Persistenzdienst und die Geschäftslogik synchron oder in strikter Reihenfolge aktualisiert werden.
Letztendlich musste ich zu den Wurzeln zurückkehren: Die Domain-Daten und die Domain-Geschäftslogik in einem Microservice zu kapseln. Dieser Ansatz vereinfacht die Wahrnehmung des Microservices als einheitliche Komponente innerhalb des Systems und vermeidet die oben genannten Probleme. Das hat jedoch seinen Preis:
- Es ist eine Standardisierung der API erforderlich, um mit der Geschäftslogik zu interagieren (insbesondere um Benutzeraktivitäten in den Geschäftsprozessen zu gewährleisten) und um die API der Plattformdienste; es ist ein sorgfältigerer Umgang mit API-Änderungen, direkter und umgekehrter Kompatibilität erforderlich;
- Es ist die Hinzufügung zusätzlicher Runtime-Bibliotheken erforderlich, um die Funktionalität der Geschäftslogik in jedem solchen Microservice sicherzustellen, und dies führt zu neuen Anforderungen an diese Bibliotheken: Leichtgewichtigkeit und minimale transitive Abhängigkeiten;
- Die Entwickler der Geschäftslogik müssen die Versionen der Bibliotheken im Auge behalten: Wenn ein Mikrodienst lange nicht aktualisiert wurde, enthält er wahrscheinlich eine veraltete Version der Bibliotheken. Dies kann ein unerwartetes Hindernis für die Implementierung neuer Funktionen darstellen und gegebenenfalls eine Migration der alten Geschäftslogik dieses Dienstes auf neue Versionen der Bibliotheken erfordern, wenn zwischen den Versionen inkompatible Änderungen vorgenommen wurden.

In einer solchen Architektur ist auch eine Schicht plattformbasierter Dienste vorhanden, aber diese Schicht bildet nicht mehr einen Container zur Ausführung der domänenspezifischen Geschäftslogik, sondern lediglich deren Umgebung, indem sie unterstützende "plattformspezifische" Funktionen bereitstellt. Eine solche Schicht ist nicht nur erforderlich, um die Leichtigkeit der domänenspezifischen Mikrodienste zu bewahren, sondern auch zur Zentralisierung des Managements.
Zum Beispiel erzeugen Benutzeraktivitäten in Geschäftsprozessen Aufgaben. Wenn Benutzer jedoch mit Aufgaben arbeiten, müssen sie Aufgaben aus allen Domänen in einer allgemeinen Liste sehen können, sodass ein entsprechender plattformbasierter Dienst zur Registrierung von Aufgaben erforderlich ist, der von der domänenspezifischen Geschäftslogik befreit ist. Die Kapselung der Geschäftslogik in einem solchen Kontext aufrechtzuerhalten, ist ziemlich schwierig, und dies ist ein weiterer Kompromiss dieser Architektur.
Integration von Geschäftsprozessen aus der Perspektive eines Anwendungsentwicklers
Wie bereits erwähnt, sollte der Anwendungsentwickler von den technischen und ingenieurtechnischen Besonderheiten der Implementierung der Interaktion mehrerer Anwendungen abstrahiert sein, um eine gute Produktivität der Entwicklung erwarten zu können.
Lassen Sie uns versuchen, eine recht anspruchsvolle Integrationsaufgabe zu lösen, die speziell für diesen Artikel konzipiert wurde. Es wird eine "Spiel"-Aufgabe mit drei Anwendungen sein, wobei jede von ihnen einen bestimmten Domänennamen definiert: "app1", "app2", "app3".
Innerhalb jeder Anwendung werden Geschäftsprozesse gestartet, die beginnen, über einen Integrationsbus "Ball zu spielen". In der Rolle des Balls fungieren Nachrichten mit dem Namen "Ball".
Spielregeln:
- Der erste Spieler – der Initiator. Er lädt andere Spieler zum Spiel ein, beginnt das Spiel und kann es jederzeit beenden;
- Andere Spieler erklären ihre Teilnahme am Spiel, "lernen" einander und den ersten Spieler kennen;
- Nach dem Empfang des Balls wählt der Spieler einen anderen teilnehmenden Spieler aus und übergibt ihm den Ball. Die Gesamtzahl der Pässe wird gezählt;
- Jeder Spieler hat "Energie", die bei jedem Pass des Balls von diesem Spieler abnimmt. Nach Erschöpfung der Energie scheidet der Spieler aus dem Spiel aus und kündigt seinen Rücktritt an;
- Wenn der Spieler allein ist, kündigt er sofort seinen Rücktritt an;
- Wenn alle Spieler ausgeschieden sind, erklärt der erste Spieler das Spiel für beendet. Wenn er vorher aus dem Spiel ausgeschieden ist, bleibt er da, um das Spiel zu beobachten und es abzuschließen.
Um diese Aufgabe zu lösen, werde ich unsere DSL für Geschäftsprozesse nutzen, die es ermöglicht, die Logik kompakt in Kotlin mit minimalem Boilerplate zu beschreiben.
Im Programm app1 wird der Geschäftsprozess des ersten Spielers (der Initiator des Spiels) ausgeführt:
class InitialPlayer
import ru.krista.bpm.ProcessInstance
import ru.krista.bpm.runtime.ProcessImpl
import ru.krista.bpm.runtime.constraint.UniqueConstraints
import ru.krista.bpm.runtime.dsl.processModel
import ru.krista.bpm.runtime.dsl.taskOperation
import ru.krista.bpm.runtime.instance.MessageSendInstance
data class PlayerInfo(val name: String, val domain: String, val id: String)
class PlayersList : ArrayList()
// Diese Klasse stellt eine Prozessinstanz dar: sie kapselt ihren internen Zustand
class InitialPlayer : ProcessImpl(initialPlayerModel) {
var playerName: String by persistent("Player1")
var energy: Int by persistent(30)
var players: PlayersList by persistent(PlayersList())
var shotCounter: Int = 0
}
// Dies ist die Deklaration des Prozessmodells: wird einmal erstellt, von allen
// Instanzen des entsprechenden Klassenprozesses verwendet
val initialPlayerModel = processModel(name = "InitialPlayer",
version = 1) {
// Laut Regeln ist der erste Spieler der Initiator des Spiels und muss der einzige sein
uniqueConstraint = UniqueConstraints.singleton
// Wir deklarieren die Aktivitäten, aus denen der Geschäftsprozess besteht
val sendNewGameSignal = signal("NewGame")
val sendStopGameSignal = signal("StopGame")
val startTask = humanTask("Start") {
taskOperation {
processCondition { players.size > 0 }
confirmation { "${players.size} Spieler haben sich verbunden. Beginnen wir?" }
}
}
val stopTask = humanTask("Stop") {
taskOperation {}
}
val waitPlayerJoin = signalWait("PlayerJoin") { signal ->
players.add(PlayerInfo(
signal.data!!,
signal.sender.domain,
signal.sender.processInstanceId))
println("... Spieler ${signal.data} tritt bei ...")
}
val waitPlayerOut = signalWait("PlayerOut") { signal ->
players.remove(PlayerInfo(
signal.data!!,
signal.sender.domain,
signal.sender.processInstanceId))
println("... Spieler ${signal.data} ist draußen ...")
}
val sendPlayerOut = signal("PlayerOut") {
signalData = { playerName }
}
val sendHandshake = messageSend("Handshake") {
messageData = { playerName }
activation = {
receiverDomain = process.players.last().domain
receiverProcessInstanceId = process.players.last().id
}
}
val throwStartBall = messageSend("Ball") {
messageData = { 1 }
activation = { selectNextPlayer() }
}
val throwBall = messageSend("Ball") {
messageData = { shotCounter + 1 }
activation = { selectNextPlayer() }
onEntry { energy -= 1 }
}
val waitBall = messageWaitData("Ball") {
shotCounter = it
}
// Jetzt konstruieren wir das Prozessdiagramm aus den deklarierten Aktivitäten
startFrom(sendNewGameSignal)
.fork("mainFork") {
next(startTask)
next(waitPlayerJoin).next(sendHandshake).next(waitPlayerJoin)
next(waitPlayerOut)
.branch("checkPlayers") {
ifTrue { players.isEmpty() }
.next(sendStopGameSignal)
.terminate()
ifElse().next(waitPlayerOut)
}
}
startTask.fork("afterStart") {
next(throwStartBall)
.branch("mainLoop") {
ifTrue { energy < 5 }.next(sendPlayerOut).next(waitBall)
ifElse().next(waitBall).next(throwBall).loop()
}
next(stopTask).next(sendStopGameSignal)
}
// Fügen wir zusätzliche Handler für Protokollierung zu den Aktivitäten hinzu
sendNewGameSignal.onExit { println("Lasst uns spielen!") }
sendStopGameSignal.onExit { println("Stop!") }
sendPlayerOut.onExit { println("$playerName: Ich bin draußen!") }
}
private fun MessageSendInstance.selectNextPlayer() {
val player = process.players.random()
receiverDomain = player.domain
receiverProcessInstanceId = player.id
println("Schritt ${process.shotCounter + 1}: " +
"${process.playerName} >>> ${player.name}")
}Neben der Umsetzung der Geschäftslogik kann der angegebene Code ein Objektmodell des Geschäftsprozesses bereitstellen, das in Form eines Diagramms visualisiert werden kann. Einen Visualisierer haben wir bisher noch nicht umgesetzt, daher mussten wir etwas Zeit mit dem Zeichnen verbringen (hier habe ich die BPMN-Notation in Bezug auf die Verwendung von Gates leicht vereinfacht, um die Konsistenz des Diagramms mit dem angegebenen Code zu verbessern):

Die Anwendung app2 wird Geschäftsprozesse eines anderen Spielers beinhalten:
class RandomPlayer
import ru.krista.bpm.ProcessInstance
import ru.krista.bpm.runtime.ProcessImpl
import ru.krista.bpm.runtime.dsl.processModel
import ru.krista.bpm.runtime.instance.MessageSendInstance
data class SpielerInfo(val name: String, val domain: String, val id: String)
class SpielerListe: ArrayList<SpielerInfo>()
class ZufallsSpieler : ProcessImpl<ZufallsSpieler>(zufallsSpielerModell) {
var spielerName: String by input(persistent = true,
defaultValue = "ZufallsSpieler")
var energie: Int by input(persistent = true, defaultValue = 30)
var spieler: SpielerListe by persistent(SPIELERLISTE())
var alleSpielerDraussen: Boolean by persistent(false)
var schussZaehler: Int = 0
val selbstSpieler: SpielerInfo
get() = SpielerInfo(spielerName, env.eventDispatcher.domainName, id)
}
val zufallsSpielerModell = processModel<ZufallsSpieler>(name = "ZufallsSpieler",
version = 1) {
val warteAufNeuesSpielSignal = signalWait<String>("NeuesSpiel")
val warteAufStopSpielSignal = signalWait<String>("StopSpiel")
val sendeSpielerBeitritt = signal<String>("SpielerBeitritt") {
signalData = { spielerName }
}
val sendeSpielerDraussen = signal<String>("SpielerDraussen") {
signalData = { spielerName }
}
val warteAufSpielerBeitritt = signalWaitCustom<String>("SpielerBeitritt") {
eventCondition = { signal ->
signal.sender.processInstanceId != process.id
&& !process.spieler.any { signal.sender.processInstanceId == it.id}
}
handler = { signal ->
spieler.add(SpielerInfo(
signal.data!!,
signal.sender.domain,
signal.sender.processInstanceId))
}
}
val warteAufSpielerDraussen = signalWait<String>("SpielerDraussen") { signal ->
spieler.remove(SpielerInfo(
signal.data!!,
signal.sender.domain,
signal.sender.processInstanceId))
alleSpielerDraussen = spieler.isEmpty()
}
val sendeHandschlag = messageSend<String>("Handschlag") {
messageData = { spielerName }
activation = {
receiverDomain = process.spieler.last().domain
receiverProcessInstanceId = process.spieler.last().id
}
}
val empfangenHandschlag = messageWait<String>("Handschlag") { message ->
if (!spieler.any { message.sender.processInstanceId == it.id}) {
spieler.add(SpielerInfo(
message.data!!,
message.sender.domain,
message.sender.processInstanceId))
}
}
val werfenBall = messageSend<Int>("Ball") {
messageData = { schussZaehler + 1 }
activation = { waehleNaechstenSpieler() }
onEntry { energie -= 1 }
}
val warteAufBall = messageWaitData<Int>("Ball") {
schussZaehler = it
}
startFrom(warteAufNeuesSpielSignal)
.fork("hauptGabel") {
next(sendeSpielerBeitritt)
.branch("hauptSchleife") {
ifTrue { energie < 5 || alleSpielerDraussen }
.next(sendeSpielerDraussen)
.next(warteAufBall)
ifElse()
.next(warteAufBall)
.next(werfenBall)
.loop()
}
next(warteAufSpielerBeitritt).next(sendeHandschlag).next(warteAufSpielerBeitritt)
next(warteAufSpielerDraussen).next(warteAufSpielerDraussen)
next(empfangenHandschlag).next(empfangenHandschlag)
next(warteAufStopSpielSignal).terminate()
}
sendeSpielerBeitritt.onExit { println("$spielerName: Ich bin hier!") }
sendeSpielerDraussen.onExit { println("$spielerName: Ich bin raus!") }
}
private fun MessageSendInstance<ZufallsSpieler, Int>.waehleNaechstenSpieler() {
val spieler = if (process.spieler.isNotEmpty())
process.spieler.random()
else
process.selbstSpieler
receiverDomain = spieler.domain
receiverProcessInstanceId = spieler.id
println("Schritt ${process.schussZaehler + 1}: " +
"${process.spielerName} >>> ${spieler.name}")
}Diagramm:

In der Anwendung app3 gestalten wir den Spieler mit einem anderen Verhalten: Anstelle einer zufälligen Auswahl des nächsten Spielers wird er nach dem Round-Robin-Algorithmus handeln:
class RoundRobinPlayer
import ru.krista.bpm.ProcessInstance
import ru.krista.bpm.runtime.ProcessImpl
import ru.krista.bpm.runtime.dsl.processModel
import ru.krista.bpm.runtime.instance.MessageSendInstance
data class PlayerInfo(val name: String, val domain: String, val id: String)
class PlayersList: ArrayList()
class RoundRobinPlayer : ProcessImpl(roundRobinPlayerModel) {
var playerName: String by input(persistent = true,
defaultValue = "RoundRobinPlayer")
var energy: Int by input(persistent = true, defaultValue = 30)
var players: PlayersList by persistent(PlayersList())
var nextPlayerIndex: Int by persistent(-1)
var allPlayersOut: Boolean by persistent(false)
var shotCounter: Int = 0
val selfPlayer: PlayerInfo
get() = PlayerInfo(playerName, env.eventDispatcher.domainName, id)
}
val roundRobinPlayerModel = processModel(
name = "RoundRobinPlayer",
version = 1) {
val waitNewGameSignal = signalWait("NewGame")
val waitStopGameSignal = signalWait("StopGame")
val sendPlayerJoin = signal("PlayerJoin") {
signalData = { playerName }
}
val sendPlayerOut = signal("PlayerOut") {
signalData = { playerName }
}
val waitPlayerJoin = signalWaitCustom("PlayerJoin") {
eventCondition = { signal ->
signal.sender.processInstanceId != process.id
&& !process.players.any { signal.sender.processInstanceId == it.id}
}
handler = { signal ->
players.add(PlayerInfo(
signal.data!!,
signal.sender.domain,
signal.sender.processInstanceId))
}
}
val waitPlayerOut = signalWait("PlayerOut") { signal ->
players.remove(PlayerInfo(
signal.data!!,
signal.sender.domain,
signal.sender.processInstanceId))
allPlayersOut = players.isEmpty()
}
val sendHandshake = messageSend("Handshake") {
messageData = { playerName }
activation = {
receiverDomain = process.players.last().domain
receiverProcessInstanceId = process.players.last().id
}
}
val receiveHandshake = messageWait("Handshake") { message ->
if (!players.any { message.sender.processInstanceId == it.id}) {
players.add(PlayerInfo(
message.data!!,
message.sender.domain,
message.sender.processInstanceId))
}
}
val throwBall = messageSend("Ball") {
messageData = { shotCounter + 1 }
activation = { selectNextPlayer() }
onEntry { energy -= 1 }
}
val waitBall = messageWaitData("Ball") {
shotCounter = it
}
startFrom(waitNewGameSignal)
.fork("mainFork") {
next(sendPlayerJoin)
.branch("mainLoop") {
ifTrue { energy < 5 || allPlayersOut }
.next(sendPlayerOut)
.next(waitBall)
ifElse()
.next(waitBall)
.next(throwBall)
.loop()
}
next(waitPlayerJoin).next(sendHandshake).next(waitPlayerJoin)
next(waitPlayerOut).next(waitPlayerOut)
next(receiveHandshake).next(receiveHandshake)
next(waitStopGameSignal).terminate()
}
sendPlayerJoin.onExit { println("$playerName: Ich bin hier!") }
sendPlayerOut.onExit { println("$playerName: Ich bin raus!") }
}
private fun MessageSendInstance.selectNextPlayer() {
var idx = process.nextPlayerIndex + 1
if (idx >= process.players.size) {
idx = 0
}
process.nextPlayerIndex = idx
val player = if (process.players.isNotEmpty())
process.players[idx]
else
process.selfPlayer
receiverDomain = player.domain
receiverProcessInstanceId = player.id
println("Schritt ${process.shotCounter + 1}: " +
"${process.playerName} >>> ${player.name}")
}Ansonsten unterscheidet sich das Verhalten des Spielers nicht von dem vorherigen, daher bleibt das Diagramm unverändert.
Jetzt wird ein Test benötigt, um das alles auszuführen. Ich führe nur den Code des Tests an, um den Artikel nicht mit Boilerplate zu überladen (tatsächlich habe ich die zuvor erstellte Testumgebung für die Integration anderer Geschäftsprozesse verwendet):
testGame()
@Test
public void testGame() throws InterruptedException {
String pl2 = startProcess(app2, "RandomPlayer", playerParams("Player2", 20));
String pl3 = startProcess(app2, "RandomPlayer", playerParams("Player3", 40));
String pl4 = startProcess(app3, "RoundRobinPlayer", playerParams("Player4", 25));
String pl5 = startProcess(app3, "RoundRobinPlayer", playerParams("Player5", 35));
String pl1 = startProcess(app1, "InitialPlayer");
// Jetzt müssen wir etwas warten, bis sich die Spieler "kennenlernen".
// Warten mit sleep ist eine schlechte Lösung, aber die einfachste.
// Machen Sie das nicht in ernsthaften Tests!
Thread.sleep(1000);
// Das Spiel starten, while die Benutzeraktivität schließen
assertTrue(closeTask(app1, pl1, "Start"));
app1.getWaiting().waitProcessFinished(pl1);
app2.getWaiting().waitProcessFinished(pl2);
app2.getWaiting().waitProcessFinished(pl3);
app3.getWaiting().waitProcessFinished(pl4);
app3.getWaiting().waitProcessFinished(pl5);
}
private Map playerParams(String name, int energy) {
Map params = new HashMap();
params.put("playerName", name);
params.put("energy", energy);
return params;
}Test starten, Log ansehen:
Konsole Ausgabe
Der Schlüssel lock://app1/process/InitialPlayer wurde gesperrt
Lass uns spielen!
Der Schlüssel lock://app1/process/InitialPlayer wurde entsperrt
Spieler2: Ich bin hier!
Spieler3: Ich bin hier!
Spieler4: Ich bin hier!
Spieler5: Ich bin hier!
... Spieler Spieler2 beitreten ...
... Spieler Spieler4 beitreten ...
... Spieler Spieler3 beitreten ...
... Spieler Spieler5 beitreten ...
Schritt 1: Spieler1 >>> Spieler3
Schritt 2: Spieler3 >>> Spieler5
Schritt 3: Spieler5 >>> Spieler3
Schritt 4: Spieler3 >>> Spieler4
Schritt 5: Spieler4 >>> Spieler3
Schritt 6: Spieler3 >>> Spieler4
Schritt 7: Spieler4 >>> Spieler5
Schritt 8: Spieler5 >>> Spieler2
Schritt 9: Spieler2 >>> Spieler5
Schritt 10: Spieler5 >>> Spieler4
Schritt 11: Spieler4 >>> Spieler2
Schritt 12: Spieler2 >>> Spieler4
Schritt 13: Spieler4 >>> Spieler1
Schritt 14: Spieler1 >>> Spieler4
Schritt 15: Spieler4 >>> Spieler3
Schritt 16: Spieler3 >>> Spieler1
Schritt 17: Spieler1 >>> Spieler2
Schritt 18: Spieler2 >>> Spieler3
Schritt 19: Spieler3 >>> Spieler1
Schritt 20: Spieler1 >>> Spieler5
Schritt 21: Spieler5 >>> Spieler1
Schritt 22: Spieler1 >>> Spieler2
Schritt 23: Spieler2 >>> Spieler4
Schritt 24: Spieler4 >>> Spieler5
Schritt 25: Spieler5 >>> Spieler3
Schritt 26: Spieler3 >>> Spieler4
Schritt 27: Spieler4 >>> Spieler2
Schritt 28: Spieler2 >>> Spieler5
Schritt 29: Spieler5 >>> Spieler2
Schritt 30: Spieler2 >>> Spieler1
Schritt 31: Spieler1 >>> Spieler3
Schritt 32: Spieler3 >>> Spieler4
Schritt 33: Spieler4 >>> Spieler1
Schritt 34: Spieler1 >>> Spieler3
Schritt 35: Spieler3 >>> Spieler4
Schritt 36: Spieler4 >>> Spieler3
Schritt 37: Spieler3 >>> Spieler2
Schritt 38: Spieler2 >>> Spieler5
Schritt 39: Spieler5 >>> Spieler4
Schritt 40: Spieler4 >>> Spieler5
Schritt 41: Spieler5 >>> Spieler1
Schritt 42: Spieler1 >>> Spieler5
Schritt 43: Spieler5 >>> Spieler3
Schritt 44: Spieler3 >>> Spieler5
Schritt 45: Spieler5 >>> Spieler2
Schritt 46: Spieler2 >>> Spieler3
Schritt 47: Spieler3 >>> Spieler2
Schritt 48: Spieler2 >>> Spieler5
Schritt 49: Spieler5 >>> Spieler4
Schritt 50: Spieler4 >>> Spieler2
Schritt 51: Spieler2 >>> Spieler5
Schritt 52: Spieler5 >>> Spieler1
Schritt 53: Spieler1 >>> Spieler5
Schritt 54: Spieler5 >>> Spieler3
Schritt 55: Spieler3 >>> Spieler5
Schritt 56: Spieler5 >>> Spieler2
Schritt 57: Spieler2 >>> Spieler1
Schritt 58: Spieler1 >>> Spieler4
Schritt 59: Spieler4 >>> Spieler1
Schritt 60: Spieler1 >>> Spieler4
Schritt 61: Spieler4 >>> Spieler3
Schritt 62: Spieler3 >>> Spieler2
Schritt 63: Spieler2 >>> Spieler5
Schritt 64: Spieler5 >>> Spieler4
Schritt 65: Spieler4 >>> Spieler5
Schritt 66: Spieler5 >>> Spieler1
Schritt 67: Spieler1 >>> Spieler5
Schritt 68: Spieler5 >>> Spieler3
Schritt 69: Spieler3 >>> Spieler4
Schritt 70: Spieler4 >>> Spieler2
Schritt 71: Spieler2 >>> Spieler5
Schritt 72: Spieler5 >>> Spieler2
Schritt 73: Spieler2 >>> Spieler1
Schritt 74: Spieler1 >>> Spieler4
Schritt 75: Spieler4 >>> Spieler1
Schritt 76: Spieler1 >>> Spieler2
Schritt 77: Spieler2 >>> Spieler5
Schritt 78: Spieler5 >>> Spieler4
Schritt 79: Spieler4 >>> Spieler3
Schritt 80: Spieler3 >>> Spieler1
Schritt 81: Spieler1 >>> Spieler5
Schritt 82: Spieler5 >>> Spieler1
Schritt 83: Spieler1 >>> Spieler4
Schritt 84: Spieler4 >>> Spieler5
Schritt 85: Spieler5 >>> Spieler3
Schritt 86: Spieler3 >>> Spieler5
Schritt 87: Spieler5 >>> Spieler2
Schritt 88: Spieler2 >>> Spieler3
Spieler2: Ich bin raus!
Schritt 89: Spieler3 >>> Spieler4
... Spieler Spieler2 ist raus ...
Schritt 90: Spieler4 >>> Spieler1
Schritt 91: Spieler1 >>> Spieler3
Schritt 92: Spieler3 >>> Spieler1
Schritt 93: Spieler1 >>> Spieler4
Schritt 94: Spieler4 >>> Spieler3
Schritt 95: Spieler3 >>> Spieler5
Schritt 96: Spieler5 >>> Spieler1
Schritt 97: Spieler1 >>> Spieler5
Schritt 98: Spieler5 >>> Spieler3
Schritt 99: Spieler3 >>> Spieler5
Schritt 100: Spieler5 >>> Spieler4
Schritt 101: Spieler4 >>> Spieler5
Spieler4: Ich bin raus!
... Spieler Spieler4 ist raus ...
Schritt 102: Spieler5 >>> Spieler1
Schritt 103: Spieler1 >>> Spieler3
Schritt 104: Spieler3 >>> Spieler1
Schritt 105: Spieler1 >>> Spieler3
Schritt 106: Spieler3 >>> Spieler5
Schritt 107: Spieler5 >>> Spieler3
Schritt 108: Spieler3 >>> Spieler1
Schritt 109: Spieler1 >>> Spieler3
Schritt 110: Spieler3 >>> Spieler5
Schritt 111: Spieler5 >>> Spieler1
Schritt 112: Spieler1 >>> Spieler3
Schritt 113: Spieler3 >>> Spieler5
Schritt 114: Spieler5 >>> Spieler3
Schritt 115: Spieler3 >>> Spieler1
Schritt 116: Spieler1 >>> Spieler3
Schritt 117: Spieler3 >>> Spieler5
Schritt 118: Spieler5 >>> Spieler1
Schritt 119: Spieler1 >>> Spieler3
Schritt 120: Spieler3 >>> Spieler5
Schritt 121: Spieler5 >>> Spieler3
Spieler5: Ich bin raus!
... Spieler Spieler5 ist raus ...
Schritt 122: Spieler3 >>> Spieler5
Schritt 123: Spieler5 >>> Spieler1
Spieler5: Ich bin raus!
Schritt 124: Spieler1 >>> Spieler3
... Spieler Spieler5 ist raus ...
Schritt 125: Spieler3 >>> Spieler1
Schritt 126: Spieler1 >>> Spieler3
Spieler1: Ich bin raus!
... Spieler Spieler1 ist raus ...
Schritt 127: Spieler3 >>> Spieler3
Spieler3: Ich bin raus!
Schritt 128: Spieler3 >>> Spieler3
... Spieler Spieler3 ist raus ...
Spieler3: Ich bin raus!
Stop!
Schritt 129: Spieler3 >>> Spieler3
Spieler3: Ich bin raus!Aus all dem kann man mehrere wichtige Schlussfolgerungen ziehen:
- Mit den notwendigen Werkzeugen können Entwickler Integrationsinteraktionen zwischen Anwendungen schaffen, ohne von der Geschäftslogik abzuweichen;
- Die Komplexität der Integrationsaufgabe, die Ingenieurkompetenzen erfordert, kann innerhalb des Frameworks verborgen werden, wenn dies von Anfang an in die Architektur des Frameworks eingebaut wird. Die Schwierigkeit der Aufgabe kann jedoch nicht verborgen werden, weshalb die Lösung einer schwierigen Aufgabe im Code entsprechend aussehen wird;
- Bei der Entwicklung der Integrationslogik muss unbedingt die eventually consistency und das Fehlen der Linearität der Zustandsänderungen aller Integrationsbeteiligten berücksichtigt werden. Dies zwingt dazu, die Logik zu komplizieren, um sie unempfindlich gegenüber der Reihenfolge des Auftretens externer Ereignisse zu machen. In unserem Beispiel muss der Spieler erst nach seiner Ankündigung des Austritts aus dem Spiel am Spiel teilnehmen: Die anderen Spieler werden weiterhin den Ball zu ihm passen, bis die Information über seinen Austritt alle Teilnehmer erreicht und verarbeitet hat. Diese Logik ergibt sich nicht aus den Spielregeln und ist eine Kompromisslösung im Rahmen der gewählten Architektur.
Lassen Sie uns nun über verschiedene Feinheiten unserer Lösung, Kompromisse und andere Aspekte sprechen.
Alle Nachrichten – in einer Warteschlange
Alle integrierten Anwendungen arbeiten mit einer einzigen Integrationsbus, die in Form eines externen Brokers, einer Warteschlange BPMQueue – für Nachrichten und einem Topic BPMTopic – für Signale (Ereignisse) dargestellt wird. Alle Nachrichten über eine Warteschlange zu leiten, ist selbst ein Kompromiss. Auf Ebene der Geschäftslogik können nun beliebig viele neue Nachrichtentypen eingeführt werden, ohne Änderungen an der Systemstruktur vorzunehmen. Dies ist eine erhebliche Vereinfachung, birgt jedoch bestimmte Risiken, die uns im Kontext unserer typischen Aufgaben nicht so bedeutend erschienen.

Es gibt jedoch einen Punkt, der zu beachten ist: Jede Anwendung filtert ihre 'eigenen' Nachrichten bereits beim Eingang aus der Warteschlange nach ihrem Domänennamen. Der Domänenname kann auch in Signalen angegeben werden, wenn die 'Sichtbarkeit' des Signals auf eine einzige Anwendung beschränkt werden muss. Dies soll die Durchsatzrate des Busses erhöhen, aber die Geschäftslogik muss jetzt mit den Domänennamen arbeiten: für die Adressierung von Nachrichten ist dies zwingend erforderlich, für Signale wünschenswert.
Gewährleistung der Zuverlässigkeit des Integrationsbusses
Die Zuverlässigkeit setzt sich aus mehreren Aspekten zusammen:
- Der gewählte Nachrichtenbroker ist ein kritischer Bestandteil der Architektur und ein einziger Ausfallpunkt: Er sollte ausreichend ausfallsicher sein. Es sollten nur bewährte Implementierungen verwendet werden, die gut unterstützt werden und eine große Community haben;
- Die hohe Verfügbarkeit des Nachrichtenbrokers muss sichergestellt werden, weshalb er physisch von den integrierten Anwendungen getrennt sein sollte (die hohe Verfügbarkeit von Anwendungen mit geschäftlicher Logik ist erheblich schwieriger und kostspieliger zu gewährleisten);
- Der Broker muss 'at least once'-Zustellgarantien bieten. Dies ist eine zwingende Anforderung für die zuverlässige Funktion des Integrationsbusses. Garantien auf 'exactly once'-Basis sind nicht erforderlich: Geschäftsprozesse sind in der Regel nicht empfindlich gegenüber der wiederholten Zustellung von Nachrichten oder Ereignissen, und bei besonderen Aufgaben, in denen dies wichtig ist, ist es einfacher, zusätzliche Überprüfungen in die Geschäftslogik einzufügen, als ständig ausreichend 'teure' Garantien zu verwenden;
- Das Senden von Nachrichten und Signalen muss in eine allgemeine Transaktion mit der Änderung des Status von Geschäftsprozessen und Domänendaten eingebunden werden. Bevorzugt wäre die Verwendung des Musters , was jedoch eine zusätzliche Tabelle in der Datenbank und einen Relayer erfordert. In JEE-Anwendungen kann dieser Aspekt mit einem lokalen JTA-Manager vereinfacht werden, jedoch muss die Verbindung zum gewählten Broker in der Lage sein, im ;
- Modus zu arbeiten. Die Handler für eingehende Nachrichten und Ereignisse müssen ebenfalls innerhalb der Transaktion zur Änderung des Status des Geschäftsprozesses arbeiten: Wenn eine solche Transaktion zurückgerollt wird, muss auch der Empfang der Nachricht annulliert werden;
- Nachrichten, deren Zustellung aufgrund von Fehlern nicht möglich war, müssen in einem separaten Speicher abgelegt werden. (Dead Letter Queue). Wir haben dafür einen separaten Plattform-Mikroservice erstellt, der solche Nachrichten in seinem Speicher ablegt, sie nach Attributen indiziert (für schnelle Gruppierung und Suche) und eine API bereitstellt, um sie anzusehen, erneut an die Zieladresse zu senden und Nachrichten zu löschen. Systemadministratoren können mit diesem Dienst über ihre Weboberfläche arbeiten;
- In den Broker-Einstellungen muss die Anzahl der Wiederholungsversuche und die Verzögerungen zwischen den Lieferungen eingestellt werden, um die Wahrscheinlichkeit zu verringern, dass Nachrichten in die DLQ gelangen (die optimalen Parameter zu bestimmen ist praktisch unmöglich, aber man kann empirisch vorgehen und sie im Verlauf des Betriebs anpassen);
- Der DLQ-Speicher muss kontinuierlich überwacht werden, und das Überwachungssystem sollte die Systemadministratoren benachrichtigen, damit sie bei auftretenden unzustellbaren Nachrichten so schnell wie möglich reagieren können. Dies wird helfen, die „Schadenszone“ eines aufgetretenen Fehlers oder einer fehlerhaften Geschäftslogik zu verringern;
- Die Integrationsbus sollte unempfindlich gegenüber temporärer Abwesenheit von Anwendungen sein: Abonnements für das Topic sollten dauerhaft sein, und der Anwendungs-Domainname sollte eindeutig sein, damit während der Abwesenheit der Anwendung ihre Nachrichten aus der Warteschlange nicht von jemand anderem verarbeitet werden;
Sichere Thread-Sicherheit der Geschäftslogik
Einem einzigen Instanz des Geschäftsprozesses können gleichzeitig mehrere Nachrichten und Ereignisse zugehen, deren Verarbeitung parallel gestartet wird. Gleichzeitig muss für den Anwendungsentwickler alles einfach und thread-sicher sein.
Die Geschäftslogik des Prozesses verarbeitet jedes externe Ereignis, das auf diesen Geschäftsprozess Einfluss hat, einzeln. Solche Ereignisse können sein:
- der Start einer Instanz des Geschäftsprozesses;
- eine Benutzeraktion, die sich auf Aktivitäten innerhalb des Geschäftsprozesses bezieht;
- das Eintreffen einer Nachricht oder eines Signals, auf das die Instanz des Geschäftsprozesses abonniert ist;
- das Auslösen eines Timers, der von der Instanz des Geschäftsprozesses gesetzt wurde;
- eine Steuerungseinwirkung über die API (z. B. das Not-Aus des Prozesses).
Jedes solche Ereignis kann den Zustand eines Geschäftsprozess-Exemplars ändern: Einige Aktivitäten können beendet und andere können beginnen, die Werte der persistenten Eigenschaften können sich ändern. Das Schließen einer beliebigen Aktivität kann zur Aktivierung einer oder mehrerer nachfolgender Aktivitäten führen. Diese können ihrerseits auf das Eintreffen anderer Ereignisse warten oder, wenn sie keine zusätzlichen Daten benötigen, in derselben Transaktion abgeschlossen werden. Vor dem Schließen der Transaktion wird der neue Zustand des Geschäftsprozesses in der Datenbank gespeichert, wo er auf das Eintreffen des nächsten externen Ereignisses wartet.
Persistente Daten des Geschäftsprozesses, die in einer relationalen Datenbank gespeichert sind, stellen einen sehr praktischen Synchronisierungspunkt für die Verarbeitung dar, wenn SELECT FOR UPDATE verwendet wird. Wenn einer Transaktion es gelungen ist, den Zustand des Geschäftsprozesses aus der Datenbank zu erhalten, um ihn zu ändern, kann keine andere Transaktion gleichzeitig denselben Zustand für eine andere Änderung abrufen, und nach Abschluss der ersten Transaktion erhält die zweite garantiert bereits den geänderten Zustand.
Durch die Verwendung von pessimistischen Sperren auf der Seite der DBMS erfüllen wir alle notwendigen Anforderungen , und wir bewahren die Möglichkeit zur Skalierung der Anwendung mit der Geschäftslogik, indem wir die Anzahl der laufenden Exemplare erhöhen.
Eine weitere Problematik sind pessimistische Sperren, die uns Deadlocks drohen, daher sollte SELECT FOR UPDATE in der Tat mit einem vernünftigen Timeout versehen werden, um Deadlocks in einigen offensichtlichen Fällen der Geschäftslogik zu vermeiden.
Ein weiteres Problem ist die Synchronisierung des Starts des Geschäftsprozesses. Solange kein Exemplar des Geschäftsprozesses existiert, gibt es auch seinen Zustand nicht in der Datenbank, weshalb die beschriebene Methode nicht geeignet ist. Wenn die Einzigartigkeit des Exemplars des Geschäftsprozesses in einem bestimmten Umfang sichergestellt werden muss, ist ein gewisser Synchronisationsmechanismus erforderlich, der mit der Prozessklasse und dem entsprechenden Umfang assoziiert ist. Um dieses Problem zu lösen, verwenden wir einen anderen Mechanismus für Sperren, der es erlaubt, eine Sperre für eine beliebige Ressource zu erhalten, die durch einen Schlüssel im URI-Format über einen externen Dienst angegeben ist.
In unseren Beispielen enthält der Geschäftsprozess InitialPlayer die Erklärung
uniqueConstraint = UniqueConstraints.singletonDaher enthält das Protokoll Nachrichten über das Sperren und Freigeben des entsprechenden Schlüssels. Für andere Geschäftsprozesse gibt es solche Nachrichten nicht: uniqueConstraint ist nicht festgelegt.
Probleme bei Geschäftsprozessen mit persistentem Zustand
Manchmal hilft der vorhandene persistente Zustand nicht nur, sondern stört auch bei der Entwicklung erheblich.
Die Probleme beginnen, wenn Änderungen an der Geschäftslogik und/oder dem Modell des Geschäftsprozesses vorgenommen werden müssen. Nicht jede solche Änderung ist mit dem alten Zustand der Geschäftsprozesse kompatibel. Wenn in der Datenbank viele 'lebende' Instanzen vorhanden sind, kann das Einbringen inkompatibler Änderungen große Schwierigkeiten verursachen, mit denen wir oft bei der Verwendung von jBPM konfrontiert waren.
Je nach Tiefe der Änderungen kann man auf zwei Wegen vorgehen:
- einen neuen Typ von Geschäftsprozess erstellen, um inkompatible Änderungen am alten zu vermeiden, und ihn anstelle des alten bei der Erstellung neuer Instanzen verwenden. Alte Instanzen werden 'wie bisher' weiterarbeiten;
- den persistenten Zustand der Geschäftsprozesse bei Aktualisierung der Geschäftslogik migrieren.
Der erste Weg ist einfacher, hat aber seine Einschränkungen und Nachteile, wie zum Beispiel:
- Duplizierung der Geschäftslogik in vielen Modellen von Geschäftsprozessen, Erhöhung des Umfangs der Geschäftslogik;
- oft ist ein sofortiger Übergang zur neuen Geschäftslogik erforderlich (im Bereich der Integrationsaufgaben fast immer);
- der Entwickler weiß nicht, wann es möglich ist, veraltete Modelle zu löschen.
In der Praxis verwenden wir beide Ansätze, haben aber eine Reihe von Entscheidungen getroffen, um uns das Leben zu erleichtern:
- In der Datenbank wird der persistente Zustand des Geschäftsprozesses in einer leicht lesbaren und verarbeitbaren Form gespeichert: im JSON-Format. Dies ermöglicht Migrationen sowohl innerhalb der Anwendung als auch extern. Im äußersten Fall kann man es auch manuell anpassen (insbesondere bei der Entwicklung während des Debuggings);
- Die integrative Geschäftslogik verwendet keine Namen von Geschäftsprozessen, damit die Implementierung eines der beteiligten Prozesse jederzeit gegen eine neue mit einem neuen Namen (z. B. 'InitialPlayerV2') ersetzt werden kann. Die Bindung erfolgt über die Namen von Nachrichten und Signalen;
- Das Prozessmodell hat eine Versionsnummer, die wir erhöhen, wenn wir inkompatible Änderungen an diesem Modell vornehmen, und diese Nummer wird zusammen mit dem Zustand des Prozessinstanz gespeichert;
- Der persistente Zustand des Prozesses wird zuerst aus der Datenbank in ein praktisches Objektmodell eingelesen, mit dem das Migrationsverfahren arbeiten kann, wenn sich die Versionsnummer des Modells geändert hat;
- Das Migrationsverfahren wird neben der Geschäftslogik platziert und "faul" für jede Instanz des Geschäftsprozesses aufgerufen, wenn sie aus der Datenbank wiederhergestellt wird;
- Wenn der Zustand aller Prozessinstanzen schnell und synchron migriert werden muss, kommen klassischere Datenbankmigrationlösungen zum Einsatz, aber dabei muss mit JSON gearbeitet werden.
Brauchen wir ein weiteres Framework für Geschäftsprozesse?
Die in diesem Artikel beschriebenen Lösungen haben es uns ermöglicht, unser Leben erheblich zu vereinfachen, den Umfang der Fragen, die auf der Ebene der Anwendungsentwicklung gelöst werden können, zu erweitern und die Ideen zur Auslagerung von Geschäftslogik in Mikrodienste attraktiver zu gestalten. Dazu wurde viel Arbeit geleistet, ein sehr "leichtgewichtiges" Framework für Geschäftsprozesse entwickelt und Dienstkomponenten zur Lösung der angegebenen Probleme im Kontext eines breiten Spektrums anwendungsbezogener Aufgaben erstellt. Wir haben den Wunsch, diese Ergebnisse zu teilen, die Entwicklung allgemeiner Komponenten unter einer freien Lizenz öffentlich zugänglich zu machen. Dies erfordert bestimmte Anstrengungen und Zeit. Ein Verständnis für die Nachfrage nach solchen Lösungen könnte für uns einen zusätzlichen Anreiz darstellen. In dem vorgeschlagenen Artikel wird sehr wenig Augenmerk auf die Möglichkeiten des Frameworks selbst gelegt, aber einige davon sind aus den präsentierten Beispielen ersichtlich. Sollten wir unser Framework doch veröffentlichen, wird ihm ein eigener Artikel gewidmet. Bis dahin wären wir dankbar, wenn Sie uns ein kurzes Feedback geben und die Frage beantworten:
Nur registrierte Benutzer können an der Umfrage teilnehmen. .
Brauchen wir ein weiteres Framework für Geschäftsprozesse?
18,8%Ja, wir suchen schon lange nach etwas Ähnlichem3
12,5%Ich würde gerne mehr über Ihre Umsetzung erfahren, könnte nützlich sein2
6,2%Wir verwenden eines der bestehenden Frameworks, denken aber über einen Wechsel nach1
18,8%Wir verwenden eines der bestehenden Frameworks, sind damit zufrieden3
18,8%Wir kommen ohne Frameworks aus3
25,0%Wir schreiben unser eigenes4
16 Benutzer haben abgestimmt. 7 Benutzer haben sich enthalten.
Quelle: habr.com
