Learn Labs
14. Stream Processing

14.4 Kafka Streams by example

① "Every Kafka Streams application MUST have an application ID.

"Apache Kafka has two stream APIs — a low-level PROCESSOR API and a high-level STREAMS DSL." The DSL "allows us to define the application by defining a CHAIN OF TRANSFORMATIONS to events... Transformations can be as simple as a filter or as complex as a stream-to-stream join. The lower-level API allows us to create OUR OWN transformations."

The lifecycle:

StreamsBuilderdefines a topology — a DAGnew KafkaStreams(topology, props)an execution objectstreams.start()starts multiple threadsstreams.close()the processing will conclude

streams.start() “will start multiple threads, each applying the processing topology to events in the stream”.

Figure 14.4.1The lifecycle

4.1 Word count — map/filter + local-state aggregation

Configuration:

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount");           // ①
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");   // ②
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG,                // ③
  Serdes.String().getClass().getName());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG,
  Serdes.String().getClass().getName());

① "Every Kafka Streams application MUST have an application ID. It is used to coordinate the instances of the application and also when NAMING THE INTERNAL LOCAL STORES AND THE TOPICS related to them. THIS NAME MUST BE UNIQUE for each Kafka Streams application working with the same Kafka cluster." ② "Kafka Streams applications also use Kafka for COORDINATION." ③ Default Serdes — "If needed, we can override these defaults later when building the streams topology."

"you can also configure the producer and consumer EMBEDDED in Kafka Streams by adding any producer or consumer config to the Properties object" — i.e. everything from Ch. 3 and Ch. 4 applies.

The topology:

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> source = builder.stream("wordcount-input");    // ①

final Pattern pattern = Pattern.compile("\\W+");

KStream<String, String> counts = source.flatMapValues(value->           // ②
  Arrays.asList(pattern.split(value.toLowerCase())))
        .map((key, value) -> new KeyValue<String, String>(value, value))// ②
        .filter((key, value) -> (!value.equals("the")))                // ③
        .groupByKey()                                                  // ④
        .count().mapValues(value-> Long.toString(value)).toStream();   // ⑤⑥

counts.to("wordcount-output");                                         // ⑦

① Point at the input topic. ② "Each event is a line of words; we split it up using a regular expression into a series of individual words. Then we take each word (currently a VALUE) and PUT IT IN THE EVENT RECORD KEY so it can be used in a group-by operation." ③ "We filter out the word the, just to show how easy filtering is." ④ "we group by key, so we now have a collection of events for each unique word." ⑤ "We count how many events we have in each collection." ⑥ "The result of counting is a Long. We convert it to a String so it will be easier for humans to read." ⑦ Write back to Kafka.

KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
// usually the stream application would be running forever
Thread.sleep(5000L);
streams.close();

💡 The deployment payoff

*"we can run the entire example on our machine WITHOUT INSTALLING ANYTHING EXCEPT APACHE KAFKA. If our input topic contains multiple partitions, we can run MULTIPLE instances of the WordCount application (just run the app in several different terminal tabs), AND WE HAVE OUR FIRST KAFKA STREAMS PROCESSING CLUSTER. The instances talk to one another and coordinate the work.

ONE OF THE BIGGEST BARRIERS TO ENTRY FOR SOME STREAM PROCESSING FRAMEWORKS IS THAT LOCAL MODE IS VERY EASY TO USE, BUT THEN TO RUN A PRODUCTION CLUSTER, WE NEED TO INSTALL YARN OR MESOS, THEN INSTALL THE PROCESSING FRAMEWORK ON ALL THOSE MACHINES, AND THEN LEARN HOW TO SUBMIT OUR APP TO THE CLUSTER. WITH KAFKA'S STREAMS API, WE JUST START MULTIPLE INSTANCES OF OUR APP — AND WE HAVE A CLUSTER. THE EXACT SAME APP IS RUNNING ON OUR DEVELOPMENT MACHINE AND IN PRODUCTION."*

(This is Ch. 1's "APIs and libraries, not a structured runtime like YARN" design stance, cashed in.)

4.2 Stock market statistics — windowed aggregation

The goal: from a stream of trades (ticker, ask price, ask size), produce per-five-second-window: best (minimum) ask price, number of trades, and average ask price — "All statistics will be updated every second."

Serdes are the main config difference:

props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, TradeSerde.class.getName());
static public final class TradeSerde extends WrapperSerde<Trade> {
    public TradeSerde() {
        super(new JsonSerializer<Trade>(), new JsonDeserializer<Trade>(Trade.class));
    }
}

"we used the Gson library from Google to generate a JSON serializer and deserializer from our Java object. Then we created a small wrapper..."

💡 "Nothing fancy, but remember to provide a Serde object for EVERY object you want to store in Kafka — INPUT, OUTPUT, AND, IN SOME CASES, INTERMEDIATE RESULTS. To make this easier, we recommend generating these Serdes through a library like Gson, Avro, Protobuf, or something similar."

The topology:

KStream<Windowed<String>, TradeStats> stats = source
    .groupByKey()                                                       // ①
    .windowedBy(TimeWindows.of(Duration.ofMillis(windowSize))           // ②
                           .advanceBy(Duration.ofSeconds(1)))
    .aggregate(                                                         // ③
        () -> new TradeStats(),                                         // ③
        (k, v, tradestats) -> tradestats.add(v),                        // ④
        Materialized.<String, TradeStats, WindowStore<Bytes, byte[]>>
            as("trade-aggregates")                                      // ⑤
           .withValueSerde(new TradeStatsSerde()))                      // ⑥
    .toStream()                                                         // ⑦
    .mapValues((trade) -> trade.computeAvgPrice());                     // ⑧

stats.to("stockstats-output",                                           // ⑨
    Produced.keySerde(
      WindowedSerdes.timeWindowedSerdeFrom(String.class, windowSize)));

① ⚠️ The groupByKey() misnomer:

"Despite its name, this operation DOES NOT DO ANY GROUPING. Rather, it ENSURES THAT THE STREAM OF EVENTS IS PARTITIONED BASED ON THE RECORD KEY. Since we wrote the data into a topic with a key and didn't modify the key before calling groupByKey(), the data is still partitioned by its key — SO THIS METHOD DOES NOTHING IN THIS CASE."

② The window: 5 seconds, advancing every second → a hopping window (advance < size), so windows overlap.

③④ The aggregation: "The aggregate method will split the stream into OVERLAPPING WINDOWS (a five-second window every second) and then apply an aggregate method on all the events in the window." First parameter = an initializer producing the result object; second = the aggregator (add updates minimum price, trade count, total prices).

⑤⑥ The state store: "windowing aggregation requires maintaining a state and a LOCAL STORE. The last parameter is the configuration of the state store. Materialized is the store configuration object, and we configure the store name as trade-aggregates. This can be ANY unique name." Plus a Serde for the aggregation result.

⑦ Table → stream: "The result of the aggregation is A TABLE with THE TICKER AND THE TIME WINDOW AS THE PRIMARY KEY and the aggregation result as the value. We are turning the table back into a stream of events."

⑧ Deriving the average: "right now the aggregation results include the SUM of prices and NUMBER of trades. We go over these records and use the existing statistics to calculate average price." — a nice illustration: you aggregate sum + count, not average, because average isn't associative.

⑨ The windowed Serde: "Since the results are part of a windowing operation, we create a WindowedSerde that stores the result in a windowed data format that INCLUDES THE WINDOW TIMESTAMP. The window size is passed as part of the Serde, EVEN THOUGH IT ISN'T USED IN THE SERIALIZATION (DESERIALIZATION REQUIRES THE WINDOW SIZE, BECAUSE ONLY THE START TIME OF THE WINDOW IS STORED IN THE OUTPUT TOPIC)."

💡 "One thing to notice is HOW LITTLE WORK WAS NEEDED TO MAINTAIN THE LOCAL STATE of the aggregation — JUST PROVIDE A SERDE AND NAME THE STATE STORE. Yet this application will SCALE TO MULTIPLE INSTANCES and AUTOMATICALLY RECOVER FROM A FAILURE of each instance by shifting processing of some partitions to one of the surviving instances."

4.3 ClickStream enrichment — both join types in one topology

The business goal: "join all three streams to get a 360-degree view into each user activity. What did the users search for? What did they click as a result? Did they change their 'interests' in their user profile? ... Product recommendations are often based on this kind of information — the user searched for bikes, clicked on links for 'Trek,' and is interested in travel, so we can advertise bikes from Trek, helmets, and bike tours to exotic locations like Nebraska."

KStream<Integer, PageView> views =                                      // ①
    builder.stream(Constants.PAGE_VIEW_TOPIC,
      Consumed.with(Serdes.Integer(), new PageViewSerde()));

KStream<Integer, Search> searches =                                     // ①
    builder.stream(Constants.SEARCH_TOPIC,
      Consumed.with(Serdes.Integer(), new SearchSerde()));

KTable<Integer, UserProfile> profiles =                                 // ②
    builder.table(Constants.USER_PROFILE_TOPIC,
      Consumed.with(Serdes.Integer(), new ProfileSerde()));

// ── STREAM-TABLE JOIN ────────────────────────────────────────────────────
KStream<Integer, UserActivity> viewsWithProfile = views.leftJoin(profiles, // ③
                (page, profile) -> {                                    // ④
                    if (profile != null)
                        return new UserActivity(
                          profile.getUserID(), profile.getUserName(),
                          profile.getZipcode(), profile.getInterests(),
                          "", page.getPage());
                    else
                       return new UserActivity(-1, "", "", null, "", page.getPage());
                    });

// ── STREAM-STREAM (WINDOWED) JOIN ───────────────────────────────────────
KStream<Integer, UserActivity> userActivityKStream =
    viewsWithProfile.leftJoin(searches,                                 // ⑤
      (userActivity, search) -> {                                       // ⑥
          if (search != null)
              userActivity.updateSearch(search.getSearchTerms());
          else
              userActivity.updateSearch("");
          return userActivity;
      },
      JoinWindows.of(Duration.ofSeconds(1)).before(Duration.ofSeconds(0)), // ⑦
      StreamJoined.with(Serdes.Integer(),                               // ⑧
                        new UserActivitySerde(),
                        new SearchSerde()));

② "We also define a KTable for the user profiles. A KTable is A MATERIALIZED STORE THAT IS UPDATED THROUGH A STREAM OF CHANGES."

③ "In a stream-table join, EACH EVENT IN THE STREAM RECEIVES INFORMATION FROM THE CACHED COPY of the profile table. We are doing a left-join, so clicks without a known user WILL BE PRESERVED."

④ 💡 "Unlike in databases, WE GET TO DECIDE HOW TO COMBINE THE TWO VALUES INTO ONE RESULT."

⑦ The interesting part:

"a stream-to-stream join is a join with A TIME WINDOW. Joining ALL clicks and searches for each user DOESN'T MAKE MUCH SENSE — we want to join each search with clicks THAT ARE RELATED TO IT, that is, clicks that occurred a short period of time AFTER the search. So we define a join window of one second. We invoke of to create a window of one second BEFORE AND AFTER each search, and then we call before with A ZERO-SECONDS INTERVAL to make sure we ONLY JOIN CLICKS THAT HAPPEN ONE SECOND AFTER EACH SEARCH AND NOT BEFORE."

searcht = 0JoinWindows.of(1s)[−1s, +1s] — symmetric.before(0s)[0s, +1s] — only after−1ssearch (t = 0)+1s

► because a click caused by a search must come after the search. Causality is expressed as window asymmetry.

Figure 14.4.24.3 ClickStream enrichment — both join types in one topology

⑧ "the Serde of the join result... a Serde for the key that both sides of the join have in common and the Serde for both values that will be included in the result."

The two patterns, summarized:

"One joins a stream with a table to ENRICH all streaming events with information in the table. This is similar to joining a FACT TABLE with a DIMENSION when running queries on a data warehouse. The second example joins two streams based on a TIME WINDOW. THIS OPERATION IS UNIQUE TO STREAM PROCESSING."


On this page