Stateless operations (filter, mapValues, flatMapValues, map, selectKey, split, merge, peek) see one record at a time and need no store. BookNest's application parses each order, drops cancelled ones, and turns the rest into one record per order line:
package booknest;
import com.fasterxml.jackson.databind.*;
import java.time.*;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
/** Revenue per genre from BookNest's orders, joined with the catalog topic of Section 6.8.7. */
public final class GenreRevenue {
static final ObjectMapper JSON = new ObjectMapper();
public static Topology build() {
StreamsBuilder builder = new StreamsBuilder();
GlobalKTable<String, String> catalog = builder.globalTable("booknest.catalog",
Consumed.with(Serdes.String(), Serdes.String())); // book id -> JSON
KGroupedStream<String, Double> byGenre = builder
.stream("booknest.orders", Consumed.with(Serdes.String(), Serdes.String())
.withTimestampExtractor((rec, prev) -> // event time
Instant.parse(read((String) rec.value()).get("order_ts").asText()).toEpochMilli()))
.mapValues(GenreRevenue::read) // stateless steps
.filter((orderId, order) -> !order.get("status").asText().equals("cancelled"))
.flatMapValues(order -> order.get("items")) // one record per line
.map((orderId, item) -> KeyValue.pair(item.get("book_id").asText(),
item.get("qty").asInt() * item.get("unit_price").asDouble())) // key: book id
.join(catalog, (bookId, revenue) -> bookId, // stream-table join
(revenue, book) -> KeyValue.pair(read(book).get("genre").asText(), revenue))
.map((bookId, genreAndRevenue) -> genreAndRevenue) // key: genre
.repartition(Repartitioned.with(Serdes.String(), Serdes.Double()).withName("by-genre"))
.groupByKey(Grouped.with(Serdes.String(), Serdes.Double())); // no second topic
byGenre.reduce(Double::sum, Materialized.as("genre-revenue")) // running total
.toStream().to("booknest.genre-revenue", Produced.with(Serdes.String(), Serdes.Double()));
byGenre.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofDays(1), Duration.ofHours(1)))
.reduce(Double::sum, Materialized.as("genre-daily")); // one total per day
return builder.build();
}
static JsonNode read(String json) {
try { return JSON.readTree(json); } catch (Exception e) { throw new RuntimeException(e); }
}
}Prefer mapValues and flatMapValues while the key stays the same: map and selectKey mark the stream for repartitioning before the next grouping. The explicit repartition names the one internal topic; without it, each aggregation would create its own.