DEV Community

Said Olano
Said Olano

Posted on

Event Driven Architectures in Java: Building Scalable, Decoupled Systems

Event Driven Architectures in Java: Building Scalable, Decoupled Systems

Introduction

Event-driven architecture (EDA) has become a cornerstone pattern for building modern, scalable distributed systems. Unlike traditional request-response architectures where components are tightly coupled through direct method calls or REST endpoints, event-driven systems communicate asynchronously through events—immutable records of something that happened in the system.

In Java ecosystems, the shift toward event-driven architectures represents a fundamental change in how we think about system design, scalability, and resilience. Whether you're building microservices that need to scale independently, real-time data processing pipelines, or complex business workflows that span multiple domains, understanding event-driven patterns is critical for modern Java architects.

This comprehensive guide explores event-driven architectures in Java: core concepts, implementation patterns, real-world examples, and best practices that will help you design systems that are not only scalable but also maintainable and evolutionarily sound.

Part 1: Core Concepts & Patterns

What Is Event-Driven Architecture?

At its core, an event-driven architecture is a design pattern where components communicate by producing, consuming, and reacting to events. An event represents a state change or significant occurrence in the system—a user registered, an order was placed, a payment was processed, a file was uploaded.

Key characteristics:

  1. Asynchronous Communication: Components don't wait for responses. The producer emits an event and continues; the consumer processes it when ready.

  2. Loose Coupling: Components depend on event schemas, not on each other's implementation details. This allows independent scaling and deployment.

  3. Event Sourcing: The immutable sequence of events becomes your source of truth, enabling perfect audit trails and event replay capabilities.

  4. Temporal Decoupling: Producers and consumers don't need to be available at the same time, enabling better system resilience.

Event-Driven Patterns

1. Event Notification Pattern

The simplest event-driven pattern. When something happens, a producer publishes an event notification to any interested consumers. This is "fire and forget"—the producer doesn't care who handles the event.

// Event class
public class UserRegisteredEvent {
    private String userId;
    private String email;
    private LocalDateTime registeredAt;

    // Constructor, getters, equals, hashCode...
}

// Producer
@Service
public class UserService {
    @Autowired
    private ApplicationEventPublisher eventPublisher;

    public void registerUser(String email) {
        User user = new User(email);
        userRepository.save(user);

        UserRegisteredEvent event = new UserRegisteredEvent(
            user.getId(), 
            email, 
            LocalDateTime.now()
        );
        eventPublisher.publishEvent(event);
    }
}

// Consumer 1: Send welcome email
@Component
public class WelcomeEmailListener {
    @EventListener
    public void onUserRegistered(UserRegisteredEvent event) {
        emailService.sendWelcomeEmail(event.getEmail());
    }
}

// Consumer 2: Initialize user profile
@Component
public class UserProfileListener {
    @EventListener
    public void onUserRegistered(UserRegisteredEvent event) {
        profileService.initializeProfile(event.getUserId());
    }
}
Enter fullscreen mode Exit fullscreen mode

2. Event Sourcing Pattern

Instead of storing only the current state, store every event that changes state. This creates a complete audit trail and enables powerful capabilities like temporal queries and event replay.

// Event base class
public abstract class DomainEvent {
    private String aggregateId;
    private LocalDateTime occurredAt;

    public DomainEvent(String aggregateId) {
        this.aggregateId = aggregateId;
        this.occurredAt = LocalDateTime.now();
    }
}

// Domain events
public class AccountCreatedEvent extends DomainEvent {
    private String accountId;
    private String accountHolder;
    private BigDecimal initialBalance;

    public AccountCreatedEvent(String accountId, String accountHolder, BigDecimal initialBalance) {
        super(accountId);
        this.accountId = accountId;
        this.accountHolder = accountHolder;
        this.initialBalance = initialBalance;
    }
}

public class MoneyDepositedEvent extends DomainEvent {
    private String accountId;
    private BigDecimal amount;

    public MoneyDepositedEvent(String accountId, BigDecimal amount) {
        super(accountId);
        this.accountId = accountId;
        this.amount = amount;
    }
}

// Event store
@Repository
public class EventStore {
    @Autowired
    private EventRepository eventRepository;

    public void append(String aggregateId, DomainEvent event) {
        eventRepository.save(event);
    }

    public List<DomainEvent> getEvents(String aggregateId) {
        return eventRepository.findByAggregateIdOrderByOccurredAt(aggregateId);
    }

    public Account rebuildAggregate(String accountId) {
        List<DomainEvent> events = getEvents(accountId);
        Account account = new Account();

        for (DomainEvent event : events) {
            if (event instanceof AccountCreatedEvent) {
                AccountCreatedEvent e = (AccountCreatedEvent) event;
                account = new Account(e.getAccountId(), e.getAccountHolder(), e.getInitialBalance());
            } else if (event instanceof MoneyDepositedEvent) {
                MoneyDepositedEvent e = (MoneyDepositedEvent) event;
                account.deposit(e.getAmount());
            }
        }

        return account;
    }
}

// Aggregate root
public class Account {
    private String accountId;
    private String accountHolder;
    private BigDecimal balance;
    private List<DomainEvent> events = new ArrayList<>();

    public void deposit(BigDecimal amount) {
        this.balance = this.balance.add(amount);
        events.add(new MoneyDepositedEvent(accountId, amount));
    }

    public List<DomainEvent> getUncommittedEvents() {
        return new ArrayList<>(events);
    }
}
Enter fullscreen mode Exit fullscreen mode

3. CQRS (Command Query Responsibility Segregation)

Separates command models (writes) from query models (reads). Commands produce events that update the write model; event handlers project those events onto optimized read models.

// Command
public class DepositMoneyCommand {
    private String accountId;
    private BigDecimal amount;

    public DepositMoneyCommand(String accountId, BigDecimal amount) {
        this.accountId = accountId;
        this.amount = amount;
    }
}

// Command handler
@Service
public class AccountCommandService {
    @Autowired
    private AccountRepository accountRepository;

    @Autowired
    private EventPublisher eventPublisher;

    public void handle(DepositMoneyCommand cmd) {
        Account account = accountRepository.findById(cmd.getAccountId());
        account.deposit(cmd.getAmount());
        accountRepository.save(account);

        MoneyDepositedEvent event = new MoneyDepositedEvent(
            cmd.getAccountId(), 
            cmd.getAmount()
        );
        eventPublisher.publish(event);
    }
}

// Read model projection
@Service
public class AccountBalanceProjection {
    @Autowired
    private AccountBalanceRepository balanceRepository;

    @EventListener
    public void on(MoneyDepositedEvent event) {
        AccountBalance balance = balanceRepository.findById(event.getAccountId());
        balance.setBalance(balance.getBalance().add(event.getAmount()));
        balanceRepository.save(balance);
    }
}

// Query
@RestController
@RequestMapping("/accounts")
public class AccountQueryController {
    @Autowired
    private AccountBalanceRepository balanceRepository;

    @GetMapping("/{accountId}/balance")
    public BigDecimal getBalance(@PathVariable String accountId) {
        return balanceRepository.findById(accountId).getBalance();
    }
}
Enter fullscreen mode Exit fullscreen mode

Event Brokers vs Event Routers

Event Brokers (Mediator Pattern): A central broker receives events and routes them to consumers. Examples: RabbitMQ, Apache Kafka.

Event Routers (Choreography): Each component knows which other components to notify. No central broker—choreography is distributed.

Most production systems use brokers for their centralized visibility and management capabilities.

Part 2: Java Implementation with Spring

Spring ApplicationEvents (In-Process)

For monolithic applications or loosely coupled components within a single process, Spring's built-in ApplicationEventPublisher provides a clean event mechanism.

@SpringBootApplication
public class EventDrivenApplication {
    public static void main(String[] args) {
        SpringApplication.run(EventDrivenApplication.class, args);
    }
}

// Event
public class OrderPlacedEvent {
    private String orderId;
    private String customerId;
    private BigDecimal total;
    private LocalDateTime timestamp;

    public OrderPlacedEvent(String orderId, String customerId, BigDecimal total) {
        this.orderId = orderId;
        this.customerId = customerId;
        this.total = total;
        this.timestamp = LocalDateTime.now();
    }

    // Getters...
}

// Producer
@Service
public class OrderService {
    @Autowired
    private OrderRepository orderRepository;

    @Autowired
    private ApplicationEventPublisher eventPublisher;

    @Transactional
    public Order createOrder(OrderRequest request) {
        Order order = new Order(request.getCustomerId(), request.getItems());
        orderRepository.save(order);

        eventPublisher.publishEvent(
            new OrderPlacedEvent(order.getId(), request.getCustomerId(), order.getTotal())
        );

        return order;
    }
}

// Consumers
@Component
public class InventoryService {
    @EventListener
    @Async
    public void onOrderPlaced(OrderPlacedEvent event) {
        // Reduce inventory
        System.out.println("Reducing inventory for order: " + event.getOrderId());
    }
}

@Component
public class PaymentService {
    @EventListener
    @Async
    public void onOrderPlaced(OrderPlacedEvent event) {
        // Process payment
        System.out.println("Processing payment for order: " + event.getOrderId());
    }
}

@Component
public class NotificationService {
    @EventListener
    @Async
    public void onOrderPlaced(OrderPlacedEvent event) {
        // Send confirmation email
        System.out.println("Sending order confirmation to: " + event.getCustomerId());
    }
}
Enter fullscreen mode Exit fullscreen mode

Spring Cloud Stream & Apache Kafka

For distributed systems, Spring Cloud Stream provides a powerful abstraction over message brokers like Kafka, RabbitMQ, or Google Pub/Sub.

spring:
  cloud:
    stream:
      kafka:
        binder:
          brokers: localhost:9092
      bindings:
        orderPlaced-out-0:
          destination: orders.placed
        orderPlaced-in-0:
          destination: orders.placed
          group: inventory-service
Enter fullscreen mode Exit fullscreen mode
// Event
@Data
@AllArgsConstructor
public class OrderPlacedEvent {
    private String orderId;
    private String customerId;
    private List<OrderItem> items;
}

// Producer
@Service
public class OrderService {
    @Autowired
    private StreamBridge streamBridge;

    public void createOrder(OrderRequest request) {
        Order order = new Order(request.getCustomerId(), request.getItems());
        streamBridge.send("orderPlaced-out-0", 
            new OrderPlacedEvent(order.getId(), request.getCustomerId(), request.getItems())
        );
    }
}

// Consumer
@Service
public class InventoryService {
    @Bean
    public java.util.function.Consumer<OrderPlacedEvent> orderPlaced() {
        return event -> {
            System.out.println("Processing order: " + event.getOrderId());
        };
    }
}
Enter fullscreen mode Exit fullscreen mode

Spring Cloud Config + Hystrix for Resilience

Build fault-tolerant event consumers with circuit breakers.

@Service
public class PaymentEventConsumer {
    @Autowired
    private PaymentGateway paymentGateway;

    @HystrixCommand(
        commandProperties = {
            @HystrixProperty(name = "execution.isolation.thread.timeoutInMilliseconds", value = "5000"),
            @HystrixProperty(name = "circuitBreaker.errorThresholdPercentage", value = "50")
        },
        fallbackMethod = "fallbackProcessPayment"
    )
    public void processPayment(PaymentEvent event) {
        paymentGateway.charge(event.getAmount(), event.getCustomerId());
    }

    public void fallbackProcessPayment(PaymentEvent event) {
        System.out.println("Payment gateway unavailable. Queueing: " + event.getOrderId());
    }
}
Enter fullscreen mode Exit fullscreen mode

Part 3: Best Practices & Production Considerations

1. Event Design

  • Immutable: Events represent the past and cannot change
  • Self-describing: Include all relevant context
  • Versioning: Use versioned event types to handle schema evolution
public class OrderPlacedEventV2 {
    private String orderId;
    private String customerId;
    private BigDecimal total;
    private String eventVersion = "2.0";
    private String shippingAddress;
}
Enter fullscreen mode Exit fullscreen mode

2. Idempotency & Exactly-Once Processing

Consumers should be idempotent—processing the same event twice produces the same result.

@Service
public class InventoryService {
    @Autowired
    private InventoryRepository inventoryRepository;

    @Autowired
    private ProcessedEventRepository processedEventRepository;

    public void onOrderPlaced(OrderPlacedEvent event) {
        if (processedEventRepository.existsById(event.getOrderId())) {
            return;
        }
        inventory.decreaseStock(event.getOrderId());
        processedEventRepository.save(new ProcessedEvent(event.getOrderId()));
    }
}
Enter fullscreen mode Exit fullscreen mode

3. Dead Letter Queues

Failed event processing should be captured and investigated.

@Service
public class OrderEventConsumer {
    @Bean
    public java.util.function.Consumer<Message<OrderPlacedEvent>> processOrder() {
        return message -> {
            try {
                OrderPlacedEvent event = message.getPayload();
                inventory.processOrder(event);
            } catch (Exception e) {
                System.err.println("Failed to process order: " + e.getMessage());
                throw e;
            }
        };
    }
}
Enter fullscreen mode Exit fullscreen mode

4. Distributed Tracing

Correlate events across services using trace IDs.

@Service
public class OrderService {
    @Autowired
    private Tracer tracer;

    public void createOrder(OrderRequest request) {
        Span span = tracer.startSpan("create_order");
        try (Tracer.SpanInScope scope = tracer.withSpan(span)) {
            Order order = new Order(request.getCustomerId(), request.getItems());
            orderRepository.save(order);
            streamBridge.send("orderPlaced-out-0", 
                new OrderPlacedEvent(
                    order.getId(), 
                    request.getCustomerId(), 
                    request.getItems(),
                    span.context().traceId()
                )
            );
        } finally {
            span.finish();
        }
    }
}
Enter fullscreen mode Exit fullscreen mode

5. Monitoring & Observability

Track event processing latency, error rates, and consumer lag.

@Service
public class OrderEventConsumer {
    @Autowired
    private MeterRegistry meterRegistry;

    private Counter processedOrders = Counter.builder("orders.processed").register(meterRegistry);
    private Timer processingTime = Timer.builder("order.processing.time").register(meterRegistry);

    public void onOrderPlaced(OrderPlacedEvent event) {
        processingTime.recordCallable(() -> {
            inventory.processOrder(event);
            return null;
        });
        processedOrders.increment();
    }
}
Enter fullscreen mode Exit fullscreen mode

6. Event Ordering Guarantees

Kafka partitions guarantee ordering within a partition. Use a consistent key for related events.

@Service
public class OrderService {
    @Autowired
    private StreamBridge streamBridge;

    public void createOrder(OrderRequest request) {
        Message<OrderPlacedEvent> message = MessageBuilder
            .withPayload(new OrderPlacedEvent(...))
            .setHeader(KafkaHeaders.MESSAGE_KEY, request.getCustomerId().getBytes())
            .build();
        streamBridge.send("orderPlaced-out-0", message);
    }
}
Enter fullscreen mode Exit fullscreen mode

Part 4: Architecture Patterns & Real-World Examples

Saga Pattern for Distributed Transactions

Orchestrate long-running transactions across multiple services using choreography or orchestration.

@Service
public class OrderSaga {
    @Autowired
    private StreamBridge streamBridge;

    @EventListener
    public void onOrderPlaced(OrderPlacedEvent event) {
        streamBridge.send("inventory.reserve", new ReserveInventoryCommand(event.getOrderId()));
    }

    @EventListener
    public void onInventoryReserved(InventoryReservedEvent event) {
        streamBridge.send("payment.charge", new ChargePaymentCommand(event.getOrderId()));
    }

    @EventListener
    public void onPaymentProcessed(PaymentProcessedEvent event) {
        streamBridge.send("orders.confirm", new ConfirmOrderCommand(event.getOrderId()));
    }

    @EventListener(condition = "#event.success == false")
    public void onPaymentFailed(PaymentFailedEvent event) {
        streamBridge.send("inventory.release", new ReleaseInventoryCommand(event.getOrderId()));
    }
}
Enter fullscreen mode Exit fullscreen mode

Conclusion

Event-driven architecture represents a paradigm shift in how Java systems communicate and scale. By embracing asynchronous, event-based communication patterns, you build systems that are:

  • Scalable: Components scale independently based on event load
  • Resilient: Temporal decoupling means failures don't cascade
  • Evolvable: New consumers can react to events without modifying producers
  • Observable: Complete audit trails through event logs

Whether you're using Spring's ApplicationEventPublisher for monolithic applications, Spring Cloud Stream for distributed systems, or raw Kafka clients for maximum control, the principles remain consistent: design immutable events, ensure idempotent processing, implement dead letter queues, and maintain observability.

The journey from traditional request-response architectures to event-driven systems requires mindset shifts around asynchrony, consistency models, and testing strategies. But the reward is powerful: systems that truly scale, evolve, and adapt to changing business requirements.

Start small—perhaps with Spring ApplicationEvents in a monolith—then gradually adopt distributed patterns like Kafka as your system grows. The Java ecosystem provides excellent tooling for every stage of that journey.

Top comments (0)