Story Opening

The meeting agreed on Kafka in five minutes and spent forty on Lena’s question. The AI team wanted Avro and a schema registry. Kabir sketched something smaller on the whiteboard: a sealed interface, three event classes, one shared Gradle module that both services would depend on.

“That’s a schema too,” Lena said. “You compile it instead of registering it. So: who can change it, and what happens to a consumer that’s two versions behind?”

Kabir’s first draft of the publisher answered a different question badly. It called kafkaTemplate.send(…) right after save, inside the @Transactional method, the way he’d written it on every Java service before. Lena pointed at Part 13’s import test, the one whose transaction rolls back halfway. “Run that and count the events. The database forgot the price change. Kafka didn’t.”


Java → Kotlin: The Quick Map

Java/Spring habitKotlin on Spring Boot 4.1Note
@JsonTypeInfo + @JsonSubTypes on an abstract event classsealed interface + @Serializable + @SerialNameThe compiler plugin writes the discriminator; consumers get an exhaustive when
kafkaTemplate.send(…) inside the @Transactional methodpublishEvent(…) + @TransactionalEventListenerSent after commit, or not at all
Flux<ServerSentEvent<T>> on WebFluxFlow<ServerSentEvent<String>> from an MVC controllerMVC adapts Flow and suspend through Reactor
CompletableFuture<T> handler methodsuspend fun handler methodAsynchronous: resumes on whichever thread completed the awaited call
@FeignClient, or a RestTemplate wrapper@HttpExchange interface with suspend functionsNeeds the WebClient-backed group
Copying the MDC into every executor taskspring.reactor.context-propagation=auto, PropagationContextElementMicrometer context propagation
KafkaContainer + @DynamicPropertySourceKafkaContainer("apache/kafka-native:…") + @ServiceConnectionTestcontainers 2, KRaft

Conceptual Deep-Dive

Reactive where it pays

Part 11’s decision table put a JPA service on virtual threads, and Part 13 did exactly that. This part doesn’t change the decision: the catalog stays on Spring MVC, Tomcat and blocking JDBC. What virtual threads can’t give it is a stream (a dashboard that hears about changes as they commit) or structure (ask five suppliers at once, give each 500 ms, keep the answers that came back, and cancel the rest). Those are coroutine problems, and Spring MVC accepts suspend functions and Flow return values in the same controllers.

graph LR API["REST clients"] --> MVC["Spring MVC
virtual threads"] MVC --> SVC["ProductService
@Transactional"] SVC --> DB[("PostgreSQL")] SVC -.->|"publishEvent"| REL["ProductEventRelay
after commit"] REL -->|"key = SKU"| K[("Kafka
product-events")] REL --> SF["SharedFlow"] SF -->|"Flow → SSE"| DASH["Dashboards"] MVC -->|"suspend, WebClient"| SUP["Supplier gateway"] K --> AI["product-ai-service
(Part 15)"]

How MVC runs a suspend function explains most of this part’s surprises. Servlet controllers can return a Mono, so Spring wraps the call in one (with mono { } from kotlinx-coroutines-reactor, which is why that’s a dependency), starts the coroutine on Dispatchers.Unconfined, and completes the request asynchronously when the Mono does. Unconfined means “run on whatever thread you’re on”: until the first suspension, that’s the Tomcat virtual thread that took the request; after it, it’s whichever thread completed the awaited call. For a WebClient call, that’s a Reactor Netty event loop. A Flow return value takes the same route through Flux. (On a WebFlux server the handler starts on an event loop too, so there, blocking work must go to Dispatchers.IO from the first line.)

Events are a contract

Lena’s question has two honest answers:

Shared Kotlin module (this series)Schema registry (Avro, Protobuf, JSON Schema)
The schema isKotlin classes, compiled into both servicesA document in a registry, checked on every produce
Consumers getReal types and an exhaustive when over the sealed hierarchyGenerated classes, or generic records
Compatibility is checkedWhen a consumer upgrades the module and compilesAt publish time, against the registry’s rules
Works forJVM services that the same teams ownAny language, any team
CostCoupled releases; a careless change breaks consumers at runtimeInfrastructure, tooling, a second build step

For two Kotlin services owned by one platform team, the shared module is a reasonable choice, as long as the producer team treats it as a published API with three rules. New fields get defaults, so old messages still decode. Consumers ignore unknown fields, so new messages decode in old consumers. And a new event type is a breaking change for consumers that haven’t upgraded: an old consumer that receives "type":"product-archived" fails to decode it, so new types ship only after every consumer handles them. The day a Python team wants these events, the registry column wins, and the JSON had better contain every field explicitly; this part’s configuration makes sure it does.

After commit, not before

kafkaTemplate.send inside a transaction couples two systems that don’t share one. The database can roll back; Kafka has already accepted the record. So the catalog publishes Spring application events inside the transaction and sends them to Kafka only when it commits. That closes the opening’s bug and leaves a smaller gap: if the process dies between the commit and the send, the event is lost. The fix for that is the transactional outbox: write the event to an outbox table in the same transaction as the change, and let a separate relay (a poller, or change-data capture with Debezium) publish it and mark it sent. Delivery becomes at-least-once, so consumers must be idempotent. With acks=all and the producer’s default idempotence, an acknowledged record is replicated and survives a broker failover, so commit-to-send is the only loss window left. This part stops at after-commit publishing; Part 15’s consumer is idempotent either way, keyed by SKU.


Technical Explanation

The event contract

// Fragment of shared/product-events/.../ProductEvent.kt
// Every change to a product, as an event. Sealed, so a consumer's when is exhaustive and a new
// event type is a compile error in every consumer until it's handled.
@Serializable
sealed interface ProductEvent {
val eventId: Uuid
val sku: String
val occurredAt: Instant
}
@Serializable
@SerialName("product-created")
data class ProductCreated(
override val eventId: Uuid,
override val sku: String,
override val occurredAt: Instant,
val product: ProductSnapshot,
) : ProductEvent
@Serializable
@SerialName("product-updated")
data class ProductUpdated(
override val eventId: Uuid,
override val sku: String,
override val occurredAt: Instant,
val product: ProductSnapshot,
) : ProductEvent
@Serializable
@SerialName("product-deleted")
data class ProductDeleted(
override val eventId: Uuid,
override val sku: String,
override val occurredAt: Instant,
) : ProductEvent
// The full state after the change, so a consumer never has to call back for details.
@Serializable
data class ProductSnapshot(
val sku: String,
val name: String,
val description: String,
val category: String,
val pricePaise: Long,
val supplier: String,
val dietaryTags: List<String> = emptyList(),
)
// The same names as the JSON "type", for Kafka headers and SSE event names. Exhaustive: a new
// subtype doesn't compile until it has a name here.
val ProductEvent.type: String
get() = when (this) {
is ProductCreated -> "product-created"
is ProductUpdated -> "product-updated"
is ProductDeleted -> "product-deleted"
}

@SerialName decides the "type" value in JSON, so the Kotlin class names can change without breaking the wire format. Every event carries a full ProductSnapshot rather than a diff: a consumer that missed messages, or replays the topic from the start, can rebuild its state from any single event. kotlin.uuid.Uuid and kotlin.time.Instant serialise as strings with kotlinx.serialization’s built-in serializers, with no opt-in needed in Kotlin 2.4.20.

// Fragment of shared/product-events/.../ProductEventSerde.kt
object ProductEvents {
const val TOPIC = "product-events"
// "type" carries the @SerialName; unknown properties are ignored so producers can add fields
// without breaking consumers that haven't upgraded yet; defaults are written out, so an empty
// dietaryTags is [] on the wire rather than missing (kotlinx.serialization omits them by default).
val json = Json {
classDiscriminator = "type"
ignoreUnknownKeys = true
encodeDefaults = true
}
fun encode(event: ProductEvent): ByteArray = json.encodeToString(ProductEvent.serializer(), event).encodeToByteArray()
fun decode(bytes: ByteArray): ProductEvent = json.decodeFromString(ProductEvent.serializer(), bytes.decodeToString())
}
// Kafka serdes, referenced by class name in the services' configuration.
class ProductEventSerializer : Serializer<ProductEvent> {
override fun serialize(topic: String, data: ProductEvent?): ByteArray? = data?.let(ProductEvents::encode)
}
class ProductEventDeserializer : Deserializer<ProductEvent> {
override fun deserialize(topic: String, data: ByteArray?): ProductEvent? = data?.let(ProductEvents::decode)
}

One setting deserves a comment. kotlinx.serialization leaves properties at their default value out of the JSON, where Jackson writes everything. Without encodeDefaults = true, a product with no dietary tags has no dietaryTags field on the wire at all, which a Kotlin consumer fills back in from the default and a consumer in any other language sees as a missing field. The module depends on kafka-clients only at compile time (each service brings the version Spring Boot manages), and its tests pin the exact JSON shape, a round trip through the Kafka serdes, and an old consumer ignoring a field it doesn’t know.

Publishing after commit

ProductService publishes an event in each write method, next to the change it describes. create publishes ProductCreated the same way:

// Fragment of product/ProductService.kt
// Each write publishes an event inside its transaction; ProductEventRelay delivers it after commit.
@Transactional
fun reprice(sku: String, pricePaise: Long): Product = bySku(sku).apply {
this.pricePaise = pricePaise
events.publishEvent(ProductUpdated(Uuid.random(), sku, now(), snapshot()))
}
@Transactional
fun delete(sku: String) {
products.delete(bySku(sku))
events.publishEvent(ProductDeleted(Uuid.random(), sku, now()))
}
private fun now() = clock.instant().toKotlinInstant()
private fun Product.snapshot() =
ProductSnapshot(sku, name, description, category.name, pricePaise, supplier.name, dietaryTags.sorted())

publishEvent hands the event to Spring’s in-process event bus; nothing leaves the JVM yet. The relay takes it from there:

// Fragment of events/ProductEventRelay.kt
// Turns committed changes into Kafka records and into a live stream for the SSE endpoint.
@Component
class ProductEventRelay(private val kafka: KafkaTemplate<String, ProductEvent>) : SmartLifecycle {
private val log = LoggerFactory.getLogger(javaClass)
// Hot and without replay: a dashboard sees changes from the moment it connects (Part 11).
private val _changes = MutableSharedFlow<ProductEvent>(extraBufferCapacity = 256)
private val shutdown = MutableSharedFlow<Unit>(replay = 1)
val subscribers: StateFlow<Int> = _changes.subscriptionCount
// Every stream ends when the application stops: stop() emits a null sentinel, takeWhile ends the
// flow on it, and replay = 1 makes streams opened after stop() end at once. Without this, an
// endless SSE response keeps Tomcat's graceful shutdown waiting for its full timeout.
val changes: Flow<ProductEvent> = merge(_changes, shutdown.map { null }).takeWhile { it != null }.filterNotNull()
// AFTER_COMMIT: a rolled-back transaction publishes nothing (the remaining gap is in the post).
// Spring logs and swallows anything thrown here, so failures are caught and logged explicitly.
@TransactionalEventListener
fun relay(event: ProductEvent) {
if (!_changes.tryEmit(event)) log.warn("A dashboard fell behind; dropped {} for {}", event.type, event.sku)
val record = ProducerRecord(ProductEvents.TOPIC, event.sku, event) // key = SKU: per-product ordering
record.headers().add("event-type", event.type.encodeToByteArray())
try {
kafka.send(record).whenComplete { _, failure ->
if (failure != null) log.error("Could not publish {} for {}", event.type, event.sku, failure)
}
} catch (e: KafkaException) { // thrown synchronously, e.g. no broker metadata within max.block.ms
log.error("Could not publish {} for {}", event.type, event.sku, e)
}
}
// SmartLifecycle's default phase stops before the web server's graceful-shutdown phase.
@Volatile
private var running = false
override fun start() {
running = true
}
override fun stop() {
shutdown.tryEmit(Unit)
running = false
}
override fun isRunning() = running
}

@TransactionalEventListener defaults to AFTER_COMMIT: the listener runs once the transaction has committed, and never if it rolls back. The record’s key is the SKU, so every event for one product lands on the same partition, in order; the event-type header lets a consumer route, filter or skip types it doesn’t know without parsing JSON.

The listener runs on the request thread after the commit, and two details there are easy to get wrong. First, Spring runs after-commit listeners in a transaction-synchronisation callback that logs and swallows whatever they throw: the commit has happened, so there’s nobody to roll back, and the request still succeeds. Second, kafka.send usually returns a future at once, but when the producer has no metadata for the topic (Kafka down at startup, for example) it blocks for up to max.block.ms, 60 seconds by default, and then throws instead of failing the future. A callback on the future would never run, so the relay catches the exception itself. With Kafka stopped and max.block.ms: 5000, a reprice returned 200 after 5.7 seconds and the log said:

ERROR … ProductEventRelay : Could not publish product-updated for SHW-1008
Caused by: org.apache.kafka.common.errors.TimeoutException: Topic product-events not present in metadata after 5000 ms.

The change committed and the event is lost: the gap the outbox closes. The relay emits to the MutableSharedFlow first, so dashboards hear about the change even when Kafka is down. tryEmit never blocks; if a subscriber falls 256 events behind, the buffer is full and the event is dropped for every subscriber, with a warning.

changes ends every stream when the application stops. stop() emits into a second flow, shutdown, which merge turns into a null that takeWhile stops at; replay = 1 means a stream opened after stop() ends immediately. The relay is a SmartLifecycle so that Spring calls stop() before the web server’s graceful shutdown, which the Gotchas explain.

Streaming changes from an MVC controller

// Fragment of web/ProductStreamController.kt
// Spring MVC runs suspend functions and Flow return values by adapting them to Reactor types.
@RestController
class ProductStreamController(
private val relay: ProductEventRelay,
private val quotes: SupplierQuotes,
) {
// Server-sent events: one per committed change, for as long as the client stays connected.
@GetMapping("/api/products/changes", produces = [MediaType.TEXT_EVENT_STREAM_VALUE])
fun changes(): Flow<ServerSentEvent<String>> = relay.changes.map { event ->
ServerSentEvent.builder(ProductEvents.json.encodeToString(ProductEvent.serializer(), event))
.event(event.type)
.id(event.eventId.toString())
.build()
}
@GetMapping("/api/products/{sku}/best-quote")
suspend fun bestQuote(@PathVariable sku: String): BestQuote = quotes.bestQuote(sku)
}

changes() returns a Flow, and because it produces text/event-stream, Spring MVC writes each element as a server-sent event and keeps the response open. Each event is serialised with the same configuration as the Kafka records, so a browser and the AI service see identical JSON. One limitation is deliberate: the SharedFlow lives in one JVM, so with three catalog instances behind a load balancer, a dashboard sees only the changes committed by the instance it’s connected to. At that scale, feed the stream from a Kafka consumer instead.

A suspend controller and a suspending HTTP client

// Fragment of quotes/SupplierQuotes.kt
data class SupplierQuote(val supplier: String, val sku: String, val pricePaise: Long)
// Part 10's supplier client, now a real HTTP client: an interface, and Spring generates the rest.
// A suspend function needs the WebClient-backed proxy (see IntegrationConfiguration).
interface SupplierQuoteClient {
@GetExchange("/suppliers/{code}/quotes/{sku}")
suspend fun quote(@PathVariable code: String, @PathVariable sku: String): SupplierQuote
}
data class BestQuote(val sku: String, val best: SupplierQuote?, val answered: Int, val asked: Int)
@Service
class SupplierQuotes(
private val client: SupplierQuoteClient,
private val suppliers: SupplierRepository,
) {
suspend fun bestQuote(sku: String): BestQuote {
// JPA blocks, so it runs on IO rather than on whatever thread resumed this coroutine.
val codes = withContext(Dispatchers.IO) { suppliers.findAll().map { it.code } }
val quotes = coroutineScope {
codes.map { code -> async { quoteOrNull(code, sku) } }.awaitAll().filterNotNull()
}
return BestQuote(sku, quotes.minByOrNull { it.pricePaise }, answered = quotes.size, asked = codes.size)
}
// One supplier's answer, or null if it is slow or failing. Only HTTP failures are caught:
// cancellation must keep propagating (Part 10).
private suspend fun quoteOrNull(code: String, sku: String): SupplierQuote? =
try {
withTimeoutOrNull(500.milliseconds) { client.quote(code, sku) }
} catch (e: WebClientException) {
null
}
}
// Fragment of config/IntegrationConfiguration.kt
// The supplier client is generated from its interface. WEB_CLIENT because its function is suspend:
// an unspecified client type means RestClient, which Spring refuses to use for suspend functions.
@Configuration(proxyBeanMethods = false)
@ImportHttpServices(group = "suppliers", types = [SupplierQuoteClient::class], clientType = HttpServiceGroup.ClientType.WEB_CLIENT)
class IntegrationConfiguration {
// Created at startup by Boot's KafkaAdmin. Three partitions; the SKU key decides which.
@Bean
fun productEventsTopic(): NewTopic = TopicBuilder.name(ProductEvents.TOPIC).partitions(3).build()
}

SupplierQuoteClient is an interface, and Spring generates the implementation from the @GetExchange annotations. @ImportHttpServices registers it in a group named suppliers, configured by spring.http.serviceclient.suppliers.* properties. The group must be WebClient-backed, and javap shows why:

// javap com.shelfwise.catalog.quotes.SupplierQuoteClient
public abstract Object quote(String, String, Continuation<? super SupplierQuote>);

A suspend function is a method that takes a callback and may return before the result exists (Part 10). A blocking client could still implement it, by blocking and then returning the result, but Spring refuses to: its HTTP-interface proxy accepts suspend methods only with a reactive adapter, because a blocking call can’t be cancelled. An unspecified client type means RestClient in Spring Framework 7, in MVC and WebFlux applications alike, so startup fails with “Kotlin Coroutines are only supported with reactive implementations” until you set clientType = WEB_CLIENT. spring-boot-starter-webclient provides WebClient without turning the application into a WebFlux one.

Why not a blocking RestClient on a virtual thread, which Part 11 would recommend for a simple call? Because of the timeout. withTimeoutOrNull cancels a coroutine at its next suspension point; a blocking call has none, so a hung supplier would hold the request until its read timeout, whatever the coroutine asked for. Cancelling the WebClient call disposes of the subscription and closes the connection, so the 500 ms budget is real.

The withContext(Dispatchers.IO) around the repository call is defensive. Here the call comes before the first suspension, so it runs on the Tomcat virtual thread anyway. Move it below the fan-out and it would run on a Reactor Netty event loop, where a blocking JDBC call stalls every WebClient connection that thread serves. Dispatchers.IO keeps it off the event loop wherever it ends up, at the cost of a bounded pool (64 threads, or the number of cores if that’s more); a virtual-thread dispatcher, as in Part 11, is the alternative.

Keeping the request id across threads

// Fragment of web/RequestIdFilter.kt
// Puts a request id into SLF4J's MDC for every log line of the request, and echoes it back.
@Component
class RequestIdFilter : OncePerRequestFilter() {
init {
// Teach Micrometer's context propagation about this MDC key, so it can be carried across threads.
ContextRegistry.getInstance().registerThreadLocalAccessor(Slf4jThreadLocalAccessor(KEY))
}
override fun doFilterInternal(request: HttpServletRequest, response: HttpServletResponse, chain: FilterChain) {
val requestId = request.getHeader(HEADER) ?: UUID.randomUUID().toString()
MDC.put(KEY, requestId)
response.setHeader(HEADER, requestId)
try {
chain.doFilter(request, response)
} finally {
MDC.remove(KEY)
}
}
companion object {
const val KEY = "requestId"
const val HEADER = "X-Request-Id"
}
}

The filter puts a request id into SLF4J’s MDC, which is a ThreadLocal, and Part 10 showed what happens to ThreadLocals in coroutines. Part 10’s fix was MDCContext() from kotlinx-coroutines-slf4j, which carries the MDC and nothing else. Spring’s mechanism is Micrometer’s context-propagation library, which carries every ThreadLocal registered with its ContextRegistry: this MDC key, Micrometer’s current observation (and so the trace id), Spring Security’s context. With spring.reactor.context-propagation: auto in application.yaml, Spring adds a PropagationContextElement to each suspend handler’s coroutine, which captures the registered values when the handler starts and restores them on every thread it resumes on. The test probes three cases:

// Fragment of ContextPropagationTest.kt
@RestController
class Probe {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default) // as a bean-owned scope would be
@GetMapping("/test/mdc/handler")
suspend fun handler() = withContext(Dispatchers.IO) { delay(1); MDC.get("requestId") ?: "missing" }
@GetMapping("/test/mdc/own-scope")
suspend fun ownScope() = scope.async { delay(1); MDC.get("requestId") ?: "missing" }.await()
@GetMapping("/test/mdc/own-scope-propagated")
suspend fun ownScopePropagated() =
scope.async(PropagationContextElement()) { delay(1); MDC.get("requestId") ?: "missing" }.await()
}
private fun requestIdSeenBy(path: String) = client.get().uri(path).header("X-Request-Id", "req-42").exchange()
.expectStatus().isOk.expectBody<String>().returnResult().responseBody
@Test
fun `the handler's coroutine keeps the request id across dispatchers`() =
assertEquals("req-42", requestIdSeenBy("/test/mdc/handler"))
@Test
fun `a coroutine from another scope starts without it`() =
assertEquals("missing", requestIdSeenBy("/test/mdc/own-scope"))
@Test
fun `PropagationContextElement carries it into that scope`() =
assertEquals("req-42", requestIdSeenBy("/test/mdc/own-scope-propagated"))

Inside the handler’s coroutine, including after withContext(Dispatchers.IO), the id is there. A coroutine started in another scope, such as a bean-owned one, doesn’t inherit the handler’s context, so it starts without it; pass PropagationContextElement() explicitly when you launch into a scope of your own.


Step-by-Step Hands-On: Events You Can Watch

Code: kotlin-for-java-survivors/services/product-catalog-service and shared/product-events, tag kotlin-for-java-survivors/part-14.

Step 1 — Dependencies and configuration.

// Fragment of build.gradle.kts
// Part 14: events, a reactive HTTP client for suspend functions, and coroutine support.
implementation(project(":shared:product-events"))
implementation("org.springframework.boot:spring-boot-starter-kafka")
implementation("org.springframework.boot:spring-boot-starter-webclient")
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-reactor")
implementation("io.micrometer:context-propagation") // carries ThreadLocals (MDC, tracing) across threads
kafka:
bootstrap-servers: localhost:9092
producer:
value-serializer: com.shelfwise.events.ProductEventSerializer # kotlinx.serialization, from shared/product-events
acks: all # the default since Kafka 3.0: acknowledged means replicated, which the post relies on
properties:
max.block.ms: 5000 # send() can block on broker metadata after commit: cap how long a request waits
http:
serviceclient:
suppliers:
base-url: http://localhost:8090 # the suppliers' quote gateway
read-timeout: 2s

Next to these, application.yaml sets spring.reactor.context-propagation: auto, and the acks comment points at the outbox discussion above. IntegrationConfiguration declares the product-events topic as a NewTopic bean, which Boot’s KafkaAdmin creates at startup with three partitions.

Step 2 — Run it and watch. infra/compose.yaml gains a Kafka 4.2 node in KRaft mode:

kafka:
# One node in KRaft mode (broker and controller in one process, no ZooKeeper). The image's
# defaults already advertise PLAINTEXT://localhost:9092 for clients on the host.
image: apache/kafka:4.2.2
ports:
- "9092:9092"
Terminal window
docker compose -f infra/compose.yaml up -d postgres kafka
./gradlew :services:product-catalog-service:bootRun
# terminal 2: the live stream
curl -sN localhost:8080/api/products/changes
# terminal 3: create, reprice, delete
curl -s -X POST localhost:8080/api/products -H 'Content-Type: application/json' \
-d '{"sku":"SHW-1501","name":"Cold Brew Coffee 250 ml","category":"BEVERAGES","pricePaise":12000,"supplierCode":"SUP-004"}'
curl -s -X PUT localhost:8080/api/products/SHW-1501/price -H 'Content-Type: application/json' -d '{"pricePaise":11500}'
curl -s -X DELETE localhost:8080/api/products/SHW-1501

Terminal 2 prints one server-sent event per committed change (the second abbreviated):

id:ea23bb56-59cc-4255-904f-7d52856ad789
event:product-created
data:{"type":"product-created","eventId":"ea23bb56-59cc-4255-904f-7d52856ad789","sku":"SHW-1501","occurredAt":"2026-10-03T15:37:55.996795Z","product":{"sku":"SHW-1501","name":"Cold Brew Coffee 250 ml","description":"","category":"BEVERAGES","pricePaise":12000,"supplier":"Shri Krishna Consumer Goods","dietaryTags":[]}}
id:258bdcbf-ef81-4167-948c-fb843864ef60
event:product-updated
data:{"type":"product-updated",…,"product":{…,"pricePaise":11500,…}}
id:218572dc-9920-4422-aa99-7c016eda4834
event:product-deleted
data:{"type":"product-deleted","eventId":"218572dc-9920-4422-aa99-7c016eda4834","sku":"SHW-1501","occurredAt":"2026-10-03T15:37:56.181868Z"}

Kafka’s console consumer shows the same three records with their keys and headers:

Terminal window
docker compose -f infra/compose.yaml exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic product-events --from-beginning \
--formatter-property print.key=true --formatter-property print.headers=true
# event-type:product-created SHW-1501 {"type":"product-created",…}
# event-type:product-updated SHW-1501 {"type":"product-updated",…}
# event-type:product-deleted SHW-1501 {"type":"product-deleted",…}

Step 3 — Kafka in tests. The test configuration gains a second container:

// Fragment of TestcontainersConfiguration.kt
// PostgreSQL and Kafka for the test context. @ServiceConnection hands their addresses and
// credentials to Spring Boot, so there are no datasource or bootstrap-server properties to sync.
@TestConfiguration(proxyBeanMethods = false)
class TestcontainersConfiguration {
@Bean
@ServiceConnection
fun postgres() = PostgreSQLContainer("postgres:18-alpine")
@Bean
@ServiceConnection
fun kafka() = KafkaContainer("apache/kafka-native:4.2.2") // KRaft, GraalVM-native: starts in about a second
}

apache/kafka-native is the same broker compiled to a native image with GraalVM; Apache labels it experimental and not for production, but it starts much faster than the JVM image, which matters when several test contexts each start their own. @ServiceConnection registers a KafkaConnectionDetails bean that Boot’s Kafka auto-configuration uses; the spring.kafka.bootstrap-servers property itself still says localhost:9092, so a test that builds a raw client must ask the connection details, as ProductEventsTest does.

Step 4 — Prove after-commit.

// Fragment of ProductEventsTest.kt
@Test
fun `create, reprice and delete each publish one event, after commit`() {
catalog.create(newProduct("SHW-1401"))
// A rolled-back import publishes nothing, although it repriced SHW-1401 before failing.
assertFailsWith<IOException> { catalog.importPriceList(FailingAfter("SHW-1401,9900\n")) }
catalog.reprice("SHW-1401", 11_000)
catalog.delete("SHW-1401")
val events = eventsFor("SHW-1401", expected = 3)
assertEquals(listOf("product-created", "product-updated", "product-deleted"), events.map { it.first })
assertIs<ProductCreated>(events[0].second)
assertEquals(11_000, assertIs<ProductUpdated>(events[1].second).product.pricePaise)
assertIs<ProductDeleted>(events[2].second)
}

eventsFor reads the topic with a plain KafkaConsumer and the shared deserializer. Between creating and repricing SHW-1401, the test runs Part 13’s failing import, which repriced SHW-1401 to 9,900 before its connection dropped. The topic holds exactly three events for the SKU, with the reprice at 11,000: the rolled-back change never left the database.

Step 5 — Streams and fan-out against a fake supplier gateway. SupplierStub stands in for the suppliers’ API on the JDK’s built-in com.sun.net.httpserver.HttpServer: SUP-003 sleeps for two seconds, SUP-005 returns 503, the rest answer. A @DynamicPropertySource points the suppliers group at it.

// Fragment of StreamingAndQuotesTest.kt
@Test
fun `the best quote comes from the suppliers that answered in time`() {
val quote = client.get().uri("/api/products/SHW-1001/best-quote").exchange()
.expectStatus().isOk
.expectBody<BestQuote>().returnResult().responseBody!!
assertEquals("SUP-001", quote.best?.supplier) // 16,100 paise, the lowest of the answers
assertEquals(5, quote.asked)
assertEquals(3, quote.answered) // SUP-003 timed out, SUP-005 returned 503
}
@Test
fun `committed changes stream as server-sent events`(): Unit = runBlocking { // ': Unit': JUnit won't run a test that returns a value
val events = WebClient.create("http://localhost:$port").get().uri("/api/products/changes")
.retrieve()
.bodyToFlux(object : ParameterizedTypeReference<ServerSentEvent<String>>() {})
.asFlow()
val first = async { withTimeout(10_000) { events.first() } }
withTimeout(10_000) { relay.subscribers.first { it > 0 } } // the SSE request is now collecting
catalog.reprice("SHW-1008", 13_000)
val event = first.await()
assertEquals("product-updated", event.event())
assertTrue(""""pricePaise":13000""" in event.data()!!)
catalog.reprice("SHW-1008", 12_500) // restore the seed price
}

Five suppliers asked, three answered, and the endpoint returned in about half a second instead of two. The streaming test subscribes with WebClient as a dashboard would, waits until the relay’s subscriptionCount shows the SSE request collecting, reprices a product and receives the event. The whole catalog suite, 26 tests across PostgreSQL and Kafka containers, runs in about a minute.


Tips, Tricks & Gotchas

Gotcha — an endless Flow holds up graceful shutdown. Boot shuts Tomcat down gracefully, waiting up to 30 seconds for active requests, and an SSE stream is a request that never ends. The server only notices a disconnected client when it next writes, so even abandoned streams count. The first version of this service took 30 seconds to stop after every test run (“Shutdown phase … ends with 1 bean still running after timeout of 30000ms: [webServerGracefulShutdown]”). The relay is now a SmartLifecycle, which Spring stops before the web server’s graceful-shutdown phase, and its stop() ends every stream; shutdown takes two seconds.

Gotcha — a test written as = runBlocking { … } may not run. An expression-bodied test returns whatever runBlocking’s last expression returns, here a Product. JUnit 6 doesn’t execute test methods that return a value; it reports a discovery warning (“must not return a value. It will not be executed”) that is easy to miss in Gradle’s output, and the build stays green. Declare : Unit, use a block body or runTest, and add junit.platform.discovery.issue.severity.critical=WARNING to src/test/resources/junit-platform.properties, which turns the warning into a build failure.

Gotcha — exceptions in an after-commit listener disappear. Spring runs @TransactionalEventListener methods after the commit and only logs what they throw, so a Java habit, letting send()’s exception propagate and trusting the caller to see a 500, loses events silently: the response is 200 and the change is committed. Catch and log (or count) inside the listener, and treat the lost-event case as an outbox requirement rather than an error path.

Gotcha — @TransactionalEventListener is silent without a transaction. If publishEvent runs outside a transaction (a method someone forgot to annotate, or a call through this that bypasses the proxy), the listener simply never runs, and the event is lost without a log line. @TransactionalEventListener(fallbackExecution = true) runs it anyway; the safer fix is to make every write method transactional, as ProductService does at class level.

Tip — put the SKU in the key, and the type in a header. The key gives per-product ordering, and would let log compaction keep only the latest event per product if you enabled it. The event-type header lets consumers skip events they don’t care about without deserialising them.


Debugging and Troubleshooting

SymptomLikely causeFix
Events published for changes that were rolled backkafkaTemplate.send inside the transactionPublish application events; send in a @TransactionalEventListener
A write committed but no event was published, and nothing was loggedpublishEvent ran outside a transactionMake the method @Transactional, or fallbackExecution = true
An event never reaches Kafka although the change committedThe send failed after commit (logged by the relay), or the process died before sending (no log at all)A transactional outbox when loss is unacceptable
Writes take max.block.ms (60 s by default) when Kafka is down, then succeed without an eventsend() blocking on topic metadata after commit, then throwing into a listener whose exceptions Spring swallowsLower max.block.ms; catch and log in the listener; an outbox relay for guaranteed delivery
”Kotlin Coroutines are only supported with reactive implementations” at startupA suspend method in a RestClient-backed HTTP service groupclientType = HttpServiceGroup.ClientType.WEB_CLIENT
WebClient calls slow down under load; event-loop threads blockedBlocking code after a suspension point in a suspend handlerwithContext(Dispatchers.IO) around blocking calls
Log lines in coroutines lose the request idMDC not propagated, or the coroutine runs in another scopespring.reactor.context-propagation=auto; PropagationContextElement() for own scopes
A consumer fails with SerializationException on a new event typeProducer released a new subtype before consumers handled itRoll consumers first; skip unknown types by the event-type header
Shutdown waits 30 s, then logs “bean still running … webServerGracefulShutdown”Open SSE streamsEnd streams in a SmartLifecycle.stop()
A test is missing from the reportTest method returns a value (expression body): Unit; fail discovery warnings in junit-platform.properties

Key Takeaways

ConceptRemember
MVC + coroutinesKeep MVC and virtual threads for blocking JPA; use suspend and Flow where you need streams and fan-out
How MVC runs themWrapped in a Mono, started on Dispatchers.Unconfined: the thread after a suspension is whoever completed the call
Event contractSealed interface + @SerialName in a shared module: compile-time types, but a published API with rules
Wire formatclassDiscriminator, ignoreUnknownKeys, and encodeDefaults = true so every field is explicit
PublishingpublishEvent inside, @TransactionalEventListener sends after commit; its exceptions are swallowed, so catch them; mind max.block.ms; the outbox closes the remaining gap
SSEFlow<ServerSentEvent<String>> with text/event-stream; per instance; end streams on shutdown
HTTP clientssuspend compiles to a Continuation parameter; Spring implements it only with WebClient (clientType = WEB_CLIENT), whose cancellation reaches the socket
Contextspring.reactor.context-propagation=auto for handlers; PropagationContextElement() for your own scopes
Testsapache/kafka-native + @ServiceConnection; JDK HttpServer for fake upstreams; : Unit on runBlocking tests

Story Closing

The events went to the shared Kafka cluster on Wednesday, and the first thing the AI team built with them was a dashboard of their own: a terminal running kafka-console-consumer, scrolling price changes from the stores.

On Friday the AI team brought Kabir the first prototype of the assistant for a demo rehearsal. It answered “gluten-free pasta under ₹200?” with the millet pasta, correctly, from a catalog export pasted into its prompt. Then someone asked about returning a mixer grinder that had stopped working, and it explained, fluently and in detail, a thirty-day no-questions-asked return policy that Shelfwise had never had.

“It’s not lying,” the AI team’s lead said. “It just doesn’t know what it doesn’t know.”

Lena looked at Kabir. “Then the next service is yours.”

In Part 15, Kabir builds product-ai-service on Spring AI 2.0: a Kafka consumer that keeps a Qdrant index in step with the catalog, answers grounded in real documents, a tool for live prices, and a local model through Ollama.


This is Part 14 of a 16-part series: “Kotlin for Java Survivors: Life After Semicolons.”