Skip to main content

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