skip to content

Using StreamsBuilder, how do you create stream(), table(), and globalTable() sources, and what does each produce?

level: juniorimportance: should knowfreq 60%

answer

  1. StreamsBuilder = DSL entry; build() -> Topology
  2. stream() -> KStream, table() -> KTable, globalTable() -> GlobalKTable
  3. Consumed: serdes, TimestampExtractor, offset reset, name
  4. Materialized: store name (queryable), serdes, caching/logging
  5. table() materializes into a local state store

basics

~10 s

StreamsBuilder is the DSL entry point. builder.stream(topic) returns a KStream (event stream), builder.table(topic) returns a KTable (latest-value-per-key, materialized), and builder.globalTable(topic) returns a GlobalKTable (fully replicated). builder.build() produces the Topology.

solid answer

~40 s

`StreamsBuilder` is the high-level DSL builder. You declare sources from Kafka topics: - `builder.stream("topic")` → a `KStream<K,V>`, an append-only event stream. - `builder.table("topic")` → a `KTable<K,V>`, materialized into a local state store (you can name it via `Materialized.as("store")` to make it queryable). - `builder.globalTable("topic")` → a `GlobalKTable<K,V>`, fully replicated to every instance. You supply serdes either via `Consumed.with(keySerde, valueSerde)` per source or via default serde configs. After wiring operators you call `builder.build()` to get an immutable `Topology`, then hand it plus a `Properties` (config) to `new KafkaStreams(topology, props)`. The `Consumed` object also lets you set a `TimestampExtractor`, `Topic.Pattern`, auto-offset-reset, and a processor name. `table()` and `globalTable()` accept a `Materialized` to control the store name, store type, and serdes.

go deeper

for a junior

Map each factory method to its output type and know build() yields a Topology.

for a middle

Use Consumed and Materialized correctly for serdes, naming, and queryability.

for a senior

Explain default-serde fallbacks, TimestampExtractor per source, and when table() materializes a store.

for a principal

Use build(Properties) topology optimization to reuse source topics as changelogs and minimize internal topics across the app.

## StreamsBuilder: the DSL entry point `StreamsBuilder` is the fluent builder for the **high-level DSL**. You start by declaring **source operators** that read from Kafka topics, then chain transformations, and finally `build()` an immutable `Topology` to run. ### The three source factory methods **1. `stream(topic)` → `KStream`** ```java KStream<String, String> views = builder.stream("page-views"); ``` Produces an **append-only event stream**. Variants accept a collection of topics or a `Pattern` (regex) to subscribe to multiple topics. **2. `table(topic)` → `KTable`** ```java KTable<String, User> users = builder.table("users", Materialized.as("users-store")); ``` Produces a **changelog/upsert table**, materialized into a **local state store**. Naming the store via `Materialized.as(...)` makes it available for **Interactive Queries**. Without an explicit `Materialized`, the DSL may still materialize it internally if needed. **3. `globalTable(topic)` → `GlobalKTable`** ```java GlobalKTable<String, Rate> rates = builder.globalTable("fx-rates"); ``` Produces a **fully-replicated table**: every application instance loads ALL partitions of the topic into its own store. For small reference data used in foreign-key joins. ### Supplying serdes and read options: `Consumed` Each source method has an overload taking a `Consumed<K,V>`: ```java builder.stream("page-views", Consumed.with(Serdes.String(), Serdes.String()) .withTimestampExtractor(new MyExtractor()) .withOffsetResetPolicy(AutoOffsetReset.earliest()) .withName("views-source")); ``` `Consumed` controls key/value **serdes**, a **`TimestampExtractor`** (how event time is read), the **auto.offset.reset** policy for THIS source, and the processor **name** (which improves `describe()` output readability). ### Controlling the store: `Materialized` `table()`/`globalTable()` (and aggregations) take a `Materialized`: ```java Materialized.<String, User, KeyValueStore<Bytes, byte[]>>as("users-store") .withKeySerde(Serdes.String()) .withValueSerde(userSerde) .withCachingDisabled(); ``` This sets store name (queryability), store type, serdes, caching, and logging (changelog) options. ### Building and running ```java Topology topology = builder.build(); // immutable graph KafkaStreams streams = new KafkaStreams(topology, props); streams.start(); ``` `build()` can also take a `Properties` so DSL **topology optimizations** (e.g. `topology.optimization=all`) can be applied — notably reusing source topics as changelogs to drop redundant repartition/changelog topics. ### Edge cases / gotchas - If you don't set serdes via `Consumed`/`Materialized`, Streams falls back to `default.key.serde` / `default.value.serde`; a mismatch surfaces as a `ClassCastException`/`SerializationException` at runtime, not compile time. - `table()` from a non-compacted topic still works, but for correct restore the changelog backing should be compacted; the source topic itself should ideally be compacted for KTable semantics. - `stream(Pattern...)` lets one source consume many topics, but they must share key/value types/serdes.

  • What does Materialized.as("store-name") give you that an unnamed table does not?
    A named store is exposed for Interactive Queries (you can look up current values via KafkaStreams.store(...)), and you get explicit control over the store name, serdes, caching, and changelog logging. An unnamed table may be materialized internally but isn't directly queryable by a known name.
  • Where do serdes come from if you call builder.stream("topic") with no Consumed?
    From the default.key.serde and default.value.serde StreamsConfig properties. A mismatch with the actual data types fails at runtime with a SerializationException/ClassCastException, not at compile time.

saying these in an interview costs you the question

  • Confusing the builder method outputs (e.g. saying table() returns a KStream).
  • Claiming build() returns a running KafkaStreams — it returns an immutable Topology you then pass to KafkaStreams.
  • Forgetting that serdes default to StreamsConfig defaults when not given via Consumed/Materialized.
  • Thinking globalTable() is just an overload of table() with the same semantics.

context