The Processor API

The DSL compiles to the Processor API, usable directly for per-record store access, timers (punctuators) and custom forwarding. DeliveryTimes emits the hours from each order's placement to its delivery:

streams/src/booknest/DeliveryTimes.java: a stateful processor wired by handJava
package booknest;
import java.time.Instant;
import org.apache.kafka.common.serialization.*;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.processor.api.*;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.state.*;
/** Processor API: hours from order_placed to order_delivered, kept per order in a state store. */
public final class DeliveryTimes extends ContextualProcessor<String, String, String, String> {
  private KeyValueStore<String, Long> placedAt;
  @Override public void init(ProcessorContext<String, String> context) {
    super.init(context);
    placedAt = context.getStateStore("placed-at");
  }
  @Override public void process(Record<String, String> event) {
    var json = GenreRevenue.read(event.value());
    long ts = Instant.parse(json.get("ts").asText()).toEpochMilli();
    switch (json.get("type").asText()) {
      case "order_placed" -> placedAt.put(event.key(), ts);
      case "order_delivered" -> {
        Long placed = placedAt.delete(event.key());                 // and forget the order
        if (placed != null) context().forward(event.withValue("" + (ts - placed) / 3_600_000));
      }
      default -> { }                                          // other event types
    }
  }
  public static Topology build() {
    var s = new StringSerializer();  var d = new StringDeserializer();
    return new Topology().addSource("events", d, d, "booknest.order-events")
        .addProcessor("delivery-times", DeliveryTimes::new, "events")
        .addStateStore(Stores.keyValueStoreBuilder(Stores.persistentKeyValueStore("placed-at"),
            Serdes.String(), Serdes.Long()), "delivery-times")
        .addSink("hours", "booknest.delivery-hours", s, s, "delivery-times");
  }
}

In 40 seconds it matched 30,876 deliveries, a mean of 148.8 hours and a median of 149. Deleting delivered orders keeps the store small; a punctuator could expire orders never delivered.