Ukratko

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
  1. Zašto Event-Driven Arhitektura (EDA)?
  2. Prednosti EDA:
  3. Izazovi EDA:
  4. Apache Kafka: Srce Event-Driven Arhitekture
  5. Quarkus: Java za Cloud-Native Svijet
  6. Ključne prednosti Quarkusa za EDA:
  7. Implementacija Event-Driven sustava s Quarkus i Kafka: Praktičan Primjer
  8. 1. Postavljanje Okoline
  9. 2. Kreiranje Quarkus Projekta
  10. 3. Definiranje Događaja (Event Model)
  11. 4. Proizvođač Događaja (Order Service)
  12. 5. Potrošači Događaja (Email Service i Inventory Service)
  13. 6. Testiranje
  14. Skalabilnost i Otpornost u Praksi
  15. 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:

  1. Zabilježiti narudžbu.
  2. Poslati potvrdu e-poštom (asinkrono).
  3. 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.

Izvori i dodatno čitanje

  1. Quarkus - Kafka guide
  2. Apache Kafka Documentation
  3. MicroProfile Reactive Messaging
  4. Confluent Blog - Event-Driven Architecture
B
Uredništvo portala

BAJT

Službeni autorski profil redakcije portala BAJT. Sadržaj priprema i provjerava uredništvo portala.