Hoe miljarden gebeurtenissen per seconde te verwerken zonder de latentie te doden
Stel je dit voor: je moet banktransacties in realtime controleren op fraude tijdens het betalingsproces. Je hebt slechts enkele milliseconden om een beslissing te nemen. Apache Kafka streamt voortdurend gebeurtenissen, maar voor elk ervan moet je de klantgeschiedenis uit de database ophalen. Als je voor elk bericht een traditionele PostgreSQL of MySQL bevraagt, zal het systeem onmiddellijk bezwijken onder de belasting.
Dit is een veelvoorkomende impasse waar engineers tegenaan lopen bij het ontwerpen van systemen met hoge belasting. Een gewone cache zoals Redis helpt individuele sleutelleesbewerkingen te versnellen, maar wanneer het gaat om het bouwen van complexe bedrijfslogica en analyses bovenop een datastream, schieten de mogelijkheden tekort. Je wordt gedwongen om caches, message brokers en verwerkingsengines van derden aan elkaar te knopen. Op dit punt is het de moeite waard om naar Hazelcast te kijken.
Hazelcast combineert gedistribueerde in-memory opslag en een stream processing engine in één systeem. In plaats van een constructie te assembleren uit drie verschillende services, krijg je een uniform platform dat in staat is om gegevens te verwerken, te verrijken en te analyseren in realtime.
Wat er binnenin het platform gebeurt
Het hart van het platform is de Jet engine. Deze is verantwoordelijk voor het bouwen van dataverwerkingspijplijnen. Jet kan even goed werken met continue streams en statische datasets, zoals buckets in Amazon S3 of tabellen in een relationele database.
De prestatiecijfers zijn interessant. Een enkele Hazelcast node kan 10 miljoen gebeurtenissen per seconde aggregeren met een latentie van maximaal 10 milliseconden. Als je de servers clustert, schaalt de doorvoer op tot een miljard gebeurtenissen per seconde.
Om query's uit te voeren op datastreams hoef je niet in low-level Java API's te duiken. Het platform ondersteunt standaard SQL. Je kunt een vertrouwde SELECT schrijven tegen de inkomende datastream, deze joinen met een in-memory tabel, en het resultaat onmiddellijk routeren naar de doelservice.
Zo ziet het verbinden met externe bronnen eruit. Out-of-the-box krijg je een set connectors:
- Apache Kafka en JMS voor het werken met wachtrijen
- Hadoop en Amazon S3 voor toegang tot bestandsopslag
- Relationele databases via standaard JDBC
- Python-modellen voor het uitvoeren van machine learning direct in de pijplijn
Gedistribueerd geheugen en coördinatie
Als je stream analytics uit de vergelijking haalt, blijft Hazelcast een gedistribueerde key-value store. Data wordt verspreid over clusternodes als partities. Ontwikkelaars hebben toegang tot vertrouwde Java-structuren (IMap, IQueue, ITopic), met als enige verschil dat ze gedistribueerd zijn over het netwerk. Point lookups op basis van sleutels duren microseconden.
Voor database-operaties worden klassieke caching-patronen ondersteund: read-through, write-through en write-behind. Bij het gebruik van write-behind slaat de applicatie data uitsluitend op in Hazelcast's RAM, en het platform spoelt deze asynchroon weg naar schijf en de hoofddatabase. Als het relationele DBMS tijdelijk uitvalt, zal je applicatie requests blijven accepteren zonder storingen.
Een apart kenmerk is microservice-coördinatie. Hazelcast kan gedistribueerde locks beheren, unieke ID-sequenties uitgeven en gedeelde tellers bijhouden. Dit elimineert de noodzaak om een apart Apache ZooKeeper-cluster te deployen en te onderhouden voor alledaagse synchronisatietaken.
Hoe het project te bouwen en uit te voeren
De broncode van het project is geschreven in Java. JDK 17 of nieuwer is vereist voor het bouwen vanuit de bron. De eenvoudigste manier om het project te bouwen is door het Maven Wrapper-script te gebruiken:
git pull origin master
./mvnw clean package -DskipTests
Het volledige bouwen met alle controles kan even duren. Als je alleen snel lokale wijzigingen wilt verifiëren, geef dan de -Dquick vlag door:
./mvnw clean package -DskipTests -Dquick
Deze parameter schakelt Javadoc-generatie, Checkstyle-controles en het bouwen van secundaire modules uit.
De test situatie is interessant. De repository heeft duizenden tests verdeeld over drie profielen:
- Het standaard
./mvnw testprofiel voert snelle integratietests uit. - Het nightly
./mvnw test -P nightly-buildprofiel bevat langzame tests die niet parallel kunnen worden uitgevoerd. - Het volledige
./mvnw test -P all-testsprofiel voert sequentieel absoluut alle controles uit met behulp van het netwerk.
Sommige tests zijn afhankelijk van Docker. Als Docker niet op je machine is geïnstalleerd, zullen die tests falen. Om ze uit te schakelen, gebruik je de -Dhazelcast.disable.docker.tests parameter. Bij het aanmaken van een Pull Request voert de CI-server van het project de volledige suite uit, dus lokaal is het voldoende om alleen tests voor je eigen module uit te voeren.
Je kunt clients niet alleen in Java schrijven. De community en het bedrijf onderhouden officiële bibliotheken voor Python, Node.js, .NET, C++ en Go.
Licentie en een paar praktische nuances
De code in de repository is opgedeeld in twee delen. De core wordt gedistribueerd onder de permissieve Apache License 2.0. Echter, sommige enterprise-functies en modules zijn beschermd door de Hazelcast Community License. Deze verbiedt het gebruik van de code om betaalde managed services (Cloud Service Provider) te creëren die concurreren met het oorspronkelijke cloudproduct van het bedrijf.
Het tweede punt betreft resourcevereisten. Omdat alle hot data in RAM verblijft, moet je een aanzienlijke hoeveelheid geheugen aanschaffen om met grote volumes te werken. Ook in een Java-omgeving moet je goed letten op Garbage Collector-instellingen om pauzes te voorkomen bij het opschonen van gigabytes aan geheugen. Echter, Hazelcast-ingenieurs verlichten dit probleem met off-heap opslag, door data buiten de Java heap te verplaatsen.
Voor wie is Hazelcast interessant
Het platform presteert goed waar realtime gebeurtenisrespons belangrijk is:
- Fraudepreventie en scoring in fintech
- Verwerking van telemetrie en hoogfrequente IoT-signalen
- Prijzen en kortingen berekenen in e-commerce op het moment van de klik van een klant
- Data synchroniseren tussen gedistribueerde datacenters (WAN-replicatie)
Als je alleen een eenvoudige cache nodig hebt voor een paar endpoints, is Hazelcast overkill: een eenvoudiger Redis zou die taak prima aankunnen. Maar als je project is gegroeid tot de schaal waar stream analytics moet samenvloeien met gedistribueerd geheugen zonder constante trips naar schijfopslag, zal Hazelcast je maanden ontwikkeling besparen.
Gerelateerde projecten