Najvažnije iz članka
- Event-Driven arhitektura (EDA) dekuplira komponente sustava, omogućuje asinkronu komunikaciju i znatno poboljšava skalabilnost i otpornost na pogreške.
- Apache Kafka služi kao robustna, distribuirana streaming platforma za objavljivanje, pretplatu, pohranu i obradu događaja visoke propusnosti i niske latencije.
- Quarkus, *cloud-native* Java framework, idealan je za EDA zbog brzog podizanja, niskog memorijskog otiska i izvrsne integracije s Kafka preko SmallRye Reactive Messaging ekstenzije.
- Praktičan primjer pokazuje kako Quarkus servisi mogu djelovati kao proizvođači i potrošači Kafka događaja, demonstrirajući kako se postiže asinkrona obrada narudžbi i otpornost uz simulaciju neuspjeha.
Sadržaj članka
- Zašto Event-Driven Arhitektura (EDA)?
- Prednosti EDA:
- Izazovi EDA:
- Apache Kafka: Srce Event-Driven Arhitekture
- Quarkus: Java za Cloud-Native Svijet
- Ključne prednosti Quarkusa za EDA:
- Implementacija Event-Driven sustava s Quarkus i Kafka: Praktičan Primjer
- 1. Postavljanje Okoline
- 2. Kreiranje Quarkus Projekta
- 3. Definiranje Događaja (Event Model)
- 4. Proizvođač Događaja (Order Service)
- 5. Potrošači Događaja (Email Service i Inventory Service)
- 6. Testiranje
- Skalabilnost i Otpornost u Praksi
- Zaključak
U današnjem svijetu, gdje se zahtjevi za brzinom, skalabilnošću i otpornošću softverskih sustava neprestano povećavaju, arhitekturni obrasci koji podržavaju ove karakteristike postaju ključni. Event-driven arhitektura (EDA) nudi moćno rješenje dekupliranjem komponenti sustava i omogućavanjem asinkrone komunikacije temeljene na događajima. U ovom članku, zaronit ćemo u EDA paradigmu i pokazati kako je implementirati koristeći dva iznimno popularna i moćna alata: Apache Kafka za distribuiranu obradu događaja i Quarkus za razvoj visokoefikasnih, cloud-native aplikacija.
Zašto Event-Driven Arhitektura (EDA)?
Tradicionalne monolitne aplikacije i sinkroni request-response modeli često nailaze na ograničenja kada je riječ o skalabilnosti, otpornosti na pogreške i agilnosti razvoja. Svaka promjena u jednoj komponenti može utjecati na cijeli sustav, a kvar jedne usluge može paralizirati ostale. EDA nudi rješenje za ove probleme.
Prednosti EDA:
- Dekompozicija i dekupliranje: Servisi komuniciraju putem događaja, bez direktne međusobne ovisnosti. To omogućuje samostalnu implementaciju, deploy i skaliranje svakog servisa.
- Skalabilnost: Sustav može horizontalno skalirati dodavanjem novih potrošača (consumer) događaja. Kafka, kao srž EDA, inherentno podržava visoku propusnost i paralelno procesiranje.
- Otpornost na pogreške (Resilience): Kvar jedne komponente ne zaustavlja cijeli sustav. Događaji se mogu perzistirati (npr. u Kafki) i ponovno procesirati kad se servis oporavi. Servisi ne moraju biti dostupni u isto vrijeme.
- Asinkrona komunikacija: Operacije se ne blokiraju čekajući odgovor. To poboljšava odzivnost sustava i iskorištavanje resursa.
- Auditing i Replay: Svi događaji mogu biti pohranjeni, pružajući potpuni audit trail. Moguće je ponovno procesirati stare događaje za oporavak ili analizu.
- Fleksibilnost i Agilnost: Lakše je dodati nove funkcionalnosti dodavanjem novih potrošača događaja bez modificiranja postojećih servisa.
Izazovi EDA:
Unatoč brojnim prednostima, EDA donosi i određene izazove:
- Kompleksnost distribuiranog sustava: Debugging i praćenje protoka događaja može biti složenije.
- Eventualna konzistentnost: Podaci možda neće biti trenutno konzistentni u svim dijelovima sustava.
- Upravljanje događajima: Potrebna je pažljiva definicija događaja i njihovog životnog ciklusa. Sheme za događaje (npr. Avro, Protobuf) su često preporučljive.
- Operativna složenost: Postavljanje i održavanje distribuiranih messaging sustava poput Kafke zahtijeva specifično znanje.
Apache Kafka: Srce Event-Driven Arhitekture
Apache Kafka je distribuirana streaming platforma koja omogućuje objavljivanje, pretplatu, pohranu i obradu streamova zapisa u stvarnom vremenu. Razvijena je u LinkedInu i kasnije donirana Apache Software Foundationu. Ključne karakteristike Kafke uključuju:
- Visoka propusnost i niska latencija: Dizajnirana za obradu milijuna poruka u sekundi.
- Skalabilnost: Horizontano skalabilna, s podrškom za klastere razasute preko više servera.
- Otpornost na pogreške: Podaci su replicirani preko više brokera.
- Trajnost (Durability): Poruke se perzistiraju na disku.
- Zadržavanje poruka: Poruke se zadržavaju određeno vrijeme (ili dok ne dosegnu određenu veličinu) čak i nakon što su pročitane, omogućujući replay.
Kafka organizira događaje u topic-e, koji su logički kanal za određenu vrstu događaja. Svaki topic je podijeljen na particije, što omogućuje paralelno procesiranje. Proizvođači (producers) pišu događaje u topic-e, a potrošači (consumers) čitaju iz njih. Svaki događaj u particiji ima jedinstveni offset.
Quarkus: Java za Cloud-Native Svijet
Quarkus je cloud-native, container-first Java framework optimiziran za GraalVM i OpenJDK HotSpot. Njegov glavni cilj je omogućiti Java aplikacijama da se podignu brže i troše manje memorije, što je ključno za microservices i serverless arhitekture. Quarkus to postiže kompilacijom u native executables i compile-time optimizacijama.
Ključne prednosti Quarkusa za EDA:
- Brzo podizanje (Fast Startup Time): Mjereno u milisekundama, što je idealno za serverless i event-driven funkcije koje se moraju brzo pokrenuti.
- Niski memorijski otisak: Značajno manja potrošnja memorije u usporedbi s tradicionalnim Java frameworkovima, što smanjuje troškove infrastrukture.
- Developer Joy: Podržava standardne Java API-je (Jakarta EE, MicroProfile) uz live coding i brze iteracije.
- Integracija s Kafkom: Quarkus nudi izvrsnu podršku za Apache Kafka putem ekstenzija kao što su SmallRye Reactive Messaging i Kafka Client.
- Reaktivni pristup: SmallRye Reactive Messaging omogućuje implementaciju reaktivnih, event-driven tokova podataka na elegantan način.
Implementacija Event-Driven sustava s Quarkus i Kafka: Praktičan Primjer
Razmotrimo jednostavan primjer: sustav za obradu narudžbi. Kada korisnik kreira novu narudžbu, želimo:
- Zabilježiti narudžbu.
- Poslati potvrdu e-poštom (asinkrono).
- Ažurirati stanje inventara (asinkrono).
Ovo su idealni kandidati za event-driven pristup.
1. Postavljanje Okoline
Prvo, trebamo pokrenuti Kafka broker. To se najlakše radi s Dockerom i Docker Composeom:
version: '3'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.5.0
hostname: zookeeper
ports:
- "2181:2181"
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
broker:
image: confluentinc/cp-kafka:7.5.0
hostname: broker
depends_on:
- zookeeper
ports:
- "9092:9092"
- "9093:9093"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
Spremite ovo kao docker-compose.yml i pokrenite s docker-compose up -d.
2. Kreiranje Quarkus Projekta
Generirajte novi Quarkus projekt s Kafka i REST ekstenzijama:
mvn io.quarkus.platform:quarkus-maven-plugin:3.8.3:create \
-DprojectGroupId=hr.bajtr \
-DprojectArtifactId=order-service \
-Dextensions="resteasy-reactive,kafka,smallrye-reactive-messaging-kafka"
cd order-service
3. Definiranje Događaja (Event Model)
Kreirajte jednostavan POJO za OrderCreatedEvent:
package hr.bajtr.order.model;
import java.math.BigDecimal;
import java.time.Instant;
public class OrderCreatedEvent {
public String orderId;
public String customerId;
public BigDecimal amount;
public Instant timestamp;
public OrderCreatedEvent() {
}
public OrderCreatedEvent(String orderId, String customerId, BigDecimal amount) {
this.orderId = orderId;
this.customerId = customerId;
this.amount = amount;
this.timestamp = Instant.now();
}
// Getteri i setteri
// ... (ili Lombok za boilerpate kod)
}
4. Proizvođač Događaja (Order Service)
Servis za narudžbe će kreirati događaj i objaviti ga na Kafka topic. Koristimo @Outgoing anotaciju iz SmallRye Reactive Messaging za slanje poruka.
package hr.bajtr.order.service;
import hr.bajtr.order.model.OrderCreatedEvent;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.ws.rs.Consumes;
import jakarta.ws.rs.POST;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.Produces;
import jakarta.ws.rs.core.MediaType;
import jakarta.ws.rs.core.Response;
import org.eclipse.microprofile.reactive.messaging.Channel;
import org.eclipse.microprofile.reactive.messaging.Emitter;
import java.math.BigDecimal;
import java.util.UUID;
@Path("/orders")
@ApplicationScoped
@Produces(MediaType.APPLICATION_JSON)
@Consumes(MediaType.APPLICATION_JSON)
public class OrderResource {
@Channel("order-out")
Emitter<OrderCreatedEvent> orderEventEmitter;
// Za demonstraciju, jednostavno tijelo zahtjeva
record CreateOrderRequest(String customerId, BigDecimal amount) {}
@POST
public Response createOrder(CreateOrderRequest request) {
String orderId = UUID.randomUUID().toString();
OrderCreatedEvent event = new OrderCreatedEvent(
orderId,
request.customerId(),
request.amount()
);
orderEventEmitter.send(event)
.whenComplete((success, failure) -> {
if (failure != null) {
System.err.println("Failed to send order event: " + failure.getMessage());
// Ovdje bi se trebalo implementirati složenije rukovanje greškama
} else {
System.out.println("Order event sent successfully: " + orderId);
}
});
return Response.accepted().entity("Order " + orderId + " created and event published.").build();
}
}
Konfiguracija u application.properties za proizvođača:
kafka.bootstrap.servers=localhost:9092
mp.messaging.outgoing.order-out.connector=smallrye-kafka
mp.messaging.outgoing.order-out.topic=order-events
mp.messaging.outgoing.order-out.value.serializer=io.quarkus.kafka.client.serialization.ObjectMapperSerializer
ObjectMapperSerializer automatski serijalizira POJO u JSON.
5. Potrošači Događaja (Email Service i Inventory Service)
Kreirajmo dva odvojena Quarkus servisa (ili, za jednostavnost, komponente unutar istog servisa za ovaj primjer) koji će konzumirati događaje.
Email Notifier (Potrošač 1)
package hr.bajtr.order.consumer;
import hr.bajtr.order.model.OrderCreatedEvent;
import jakarta.enterprise.context.ApplicationScoped;
import org.eclipse.microprofile.reactive.messaging.Incoming;
import org.eclipse.microprofile.reactive.messaging.Message;
import java.util.concurrent.CompletionStage;
@ApplicationScoped
public class EmailNotifier {
@Incoming("orders-in")
public CompletionStage<Void> processOrderCreatedEvent(Message<OrderCreatedEvent> message) {
OrderCreatedEvent event = message.getPayload();
System.out.println("Email Service: Sending email for Order ID: " + event.orderId + " to customer: " + event.customerId);
// Simulacija slanja e-pošte
return message.ack(); // Potvrdi da je poruka uspješno obrađena
}
}
Inventory Updater (Potrošač 2)
package hr.bajtr.order.consumer;
import hr.bajtr.order.model.OrderCreatedEvent;
import jakarta.enterprise.context.ApplicationScoped;
import org.eclipse.microprofile.reactive.messaging.Incoming;
import org.eclipse.microprofile.reactive.messaging.Message;
import java.util.concurrent.CompletionStage;
@ApplicationScoped
public class InventoryUpdater {
@Incoming("orders-in")
public CompletionStage<Void> updateInventory(Message<OrderCreatedEvent> message) {
OrderCreatedEvent event = message.getPayload();
System.out.println("Inventory Service: Updating inventory for Order ID: " + event.orderId + ", amount: " + event.amount);
// Simulacija ažuriranja inventara u bazi podataka
if (Math.random() < 0.1) { // Simulacija povremenog neuspjeha
System.err.println("Inventory Update failed for order: " + event.orderId + " - will retry!");
return message.nack(new RuntimeException("Failed to update inventory")); // Negativna potvrda, poruka će biti ponovno dostavljena
}
return message.ack();
}
}
Konfiguracija u application.properties za potrošače:
kafka.bootstrap.servers=localhost:9092
# Konfiguracija za potrošače
mp.messaging.incoming.orders-in.connector=smallrye-kafka
mp.messaging.incoming.orders-in.topic=order-events
mp.messaging.incoming.orders-in.group.id=order-consumers-group # Grupni ID za potrošače
mp.messaging.incoming.orders-in.value.deserializer=io.quarkus.kafka.client.serialization.ObjectMapperDeserializer
Važno je napomenuti da group.id omogućuje da više instanci istog potrošača (npr. više instanci EmailNotifier-a) rade zajedno, pri čemu svaka instanca obrađuje podskup particija, što pruža skalabilnost i load balancing.
6. Testiranje
Pokrenite Quarkus aplikaciju u dev modu:
mvn quarkus:dev
Pošaljite narudžbu putem curl-a:
curl -X POST -H "Content-Type: application/json" -d '{"customerId": "user123", "amount": 99.99}' http://localhost:8080/orders
U konzoli Quarkus aplikacije trebali biste vidjeti logove od OrderResource, EmailNotifier i InventoryUpdater, što potvrđuje da su događaji uspješno objavljeni i konzumirani.
Skalabilnost i Otpornost u Praksi
Skalabilnost: Ako se poveća broj narudžbi, jednostavno možete pokrenuti više instanci Quarkus aplikacije (kao zasebne kontejnere ili procese). Kafka će automatski distribuirati poruke među potrošačima unutar iste group.id, osiguravajući paralelnu obradu i ravnomjerno opterećenje.
Otpornost: U primjeru InventoryUpdater, simulirali smo povremeni neuspjeh. Kada InventoryUpdater vrati message.nack(), Kafka će, ovisno o konfiguraciji, pokušati ponovno dostaviti poruku (nakon određenog timeout-a ili više puta). Ovo osigurava da se kritične operacije, poput ažuriranja inventara, eventualno izvrše čak i ako dođe do privremenih problema. U produkcijskim sustavima često se koristi i Dead Letter Queue (DLQ) za poruke koje se ne mogu obraditi nakon višestrukih pokušaja.
Zaključak
Event-driven arhitektura, kada se pravilno implementira s alatima poput Apache Kafke i Quarkusa, nudi moćan okvir za izgradnju visoko skalabilnih, otpornih i dekupliranih sustava. Quarkus sa svojim cloud-native pristupom i izvrsnom integracijom s reaktivnim messaging standardima, čini Javu idealnim izborom za razvoj komponenata u EDA. Iako donosi određenu kompleksnost, prednosti u smislu performansi, agilnosti i pouzdanosti daleko nadmašuju izazove, čineći je ključnom paradigmom u modernom razvoju softvera.
Pridržavajući se principa EDA i koristeći snagu Kafke i Quarkusa, developeri mogu izgraditi sustave koji su spremni odgovoriti na izazove sutrašnjice, s lakoćom prilagođavajući se promjenjivim poslovnim zahtjevima i nepredvidivim opterećenjima.
Komentari