Kafka
API specification
DomainEvent model
To emit a Domain Event we need to know the DomainEvent structure, which is represented with the next class:
public class DomainEvent<T> {
private final String name;
private final String eventId;
private final T data;
}
Where name is the event name, eventId is an unique event identifier and data is a JSON Serializable payload.
DomainEventBus interface
public interface DomainEventBus {
<T> Publisher<Void> emit(DomainEvent<T> event);
<T> Publisher<Void> emit(String domain, DomainEvent<T> event);
Publisher<Void> emit(CloudEvent event);
Publisher<Void> emit(String domain, CloudEvent event);
Publisher<Void> emit(RawMessage event);
Publisher<Void> emit(String domain, RawMessage event);
}
Enabling autoconfiguration
To send Domain Events you should enable the respecting spring boot autoconfiguration using the @EnableDomainEventBus
annotation For example:
@RequiredArgsConstructor
@EnableDomainEventBus
public class ReactiveEventsGateway {
public static final String SOME_EVENT_NAME = "some.event.name";
private final DomainEventBus domainEventBus; // Auto injected bean created by the @EnableDomainEventBus annotation
public Mono<Void> emit(Object event/*change for proper model*/) {
return Mono.from(domainEventBus.emit(new DomainEvent<>(SOME_EVENT_NAME, UUID.randomUUID().toString(), event)));
}
}
After that you can emit events from you application.
Sending a Raw Message
DomainEventBus.emit(RawMessage event) bypasses the DomainEvent / CloudEvent conventions: instead of building a
generic envelope, you hand over the broker-specific message yourself, with full control over its body, routing and
headers.
For Kafka there is no separate routing key parameter: the topic and the partitioning key travel inside the
KafkaMessage itself, through KafkaMessageProperties:
@RequiredArgsConstructor
@EnableDomainEventBus
public class ReactiveEventsGateway {
private final JsonMapper jsonMapper;
private final DomainEventBus domainEventBus;
public Mono<Void> emitRaw(Object payload/*change for proper model*/) {
KafkaMessage.KafkaMessageProperties properties = new KafkaMessage.KafkaMessageProperties();
properties.setTopic("some.event.name"); // required: this is where the record is published
properties.setKey(UUID.randomUUID().toString()); // optional: Kafka partitioning key
properties.getHeaders().put("content-type", "application/json");
var rawMessage = new KafkaMessage(jsonMapper.writeValueAsBytes(payload), properties, null);
return Mono.from(domainEventBus.emit(rawMessage));
}
}
Example
You can see a real example at samples/async/async-sender-client