Interactive Queries

A store is also a read model: interactive queries read the application's own state with no database in between. App queries genre-revenue every two seconds until five answers agree, then prints them (A BookNest Streams Application runs it):

streams/src/booknest/App.java: run a topology and query its storeJava
package booknest;
import java.util.Properties;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.errors.InvalidStateStoreException;
import org.apache.kafka.streams.state.QueryableStoreTypes;
/** Runs GenreRevenue until its totals settle and prints them; with "delivery", DeliveryTimes. */
public final class App {
  public static void main(String[] args) throws Exception {
    boolean delivery = args.length > 0 && args[0].equals("delivery");
    Properties p = new Properties();
    p.put(StreamsConfig.APPLICATION_ID_CONFIG, "booknest-" + (delivery ? "delivery" : "genres"));
    p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:33092");
    p.put(StreamsConfig.STATE_DIR_CONFIG, "state");                  // RocksDB files go here
    Topology topology = delivery ? DeliveryTimes.build() : GenreRevenue.build();
    try (var streams = new KafkaStreams(topology, p)) {
      streams.start();
      if (delivery) { Thread.sleep(40_000); return; }                 // Section 6.10.6
      String previous = "";
      for (int same = 0; same < 5; ) {                  // until 5 queries, 2 s apart, agree
        Thread.sleep(2_000);
        StringBuilder totals = new StringBuilder();
        try (var it = streams.store(StoreQueryParameters.fromNameAndType("genre-revenue",
            QueryableStoreTypes.<String, Double>keyValueStore())).all()) {
          it.forEachRemaining(kv -> totals.append(String.format("%s %.2f%n", kv.key, kv.value)));
        } catch (InvalidStateStoreException notQueryableYet) { continue; }  // still restoring
        same = !totals.isEmpty() && totals.toString().equals(previous) ? same + 1 : 0;
        previous = totals.toString();
      }
      System.out.print(previous);
    }
  }
}

A service would expose this over HTTP or gRPC. Each instance holds only its partitions' keys: with application.server set on each, queryMetadataForKey() names a key's owner so requests can be forwarded. Stores answer only while the instance is RUNNING, hence the retry.