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:
Asynchronous Communication: Components don't wait for responses. The producer emits an event and continues; the consumer processes it when ready.
Loose Coupling: Components depend on event schemas, not on each other's implementation details. This allows independent scaling and deployment.
Event Sourcing: The immutable sequence of events becomes your source of truth, enabling perfect audit trails and event replay capabilities.
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());
}
}
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);
}
}
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();
}
}
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());
}
}
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
// 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());
};
}
}
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());
}
}
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;
}
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()));
}
}
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;
}
};
}
}
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();
}
}
}
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();
}
}
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);
}
}
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()));
}
}
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)