Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Streaming joins: enrichment patterns

Batch & Streaming · Streaming Design

Streaming joins: enrichment patterns

Hardbatch-streaming-42
stream-joinsenrichmenttemporal-tablesflink

Question

How do you enrich a stream of orders with customer data that changes?

Solution

Enriching an order stream with mutable customer attributes can be achieved by joining against a Kafka compacted changelog topic, querying an external database using asynchronous I/O with caching, or holding reference data in broadcast state. If historical accuracy is required, temporal table joins correlate the order event timestamp with the exact version of the customer record valid at that moment. Each pattern balances output latency, memory footprint, and dimension freshness.

Stream-table joins and changelog topics

In this pattern, the customer dimension is ingested continuously from a database change data capture (CDC) feed published to a log-compacted Kafka topic:

SELECT
  o.order_id,
  o.order_amount,
  c.customer_name,
  c.loyalty_tier
FROM orders AS o
JOIN customers FOR SYSTEM_TIME AS OF o.order_time AS c
  ON o.customer_id = c.customer_id;

The streaming engine maintains the customer changelog in a local state store such as a Flink Table or Kafka Streams KTable.

Key architectural trade-offs of stream-table joins:

  • The streaming engine reads customer data from local RocksDB storage, eliminating external network round-trips and maximizing throughput.
  • Worker nodes must retain the full customer dimension in local storage, which increases cluster memory and disk overhead.

Async lookups and broadcast state

When storing the entire customer dimension locally is impractical, alternative lookup strategies are applied:

  • Asynchronous external lookups: The streaming operator calls an external store like Redis or Cassandra. To prevent network latency from blocking the stream, use asynchronous I/O with connection pools and local in-memory LRU caching. However, external calls introduce failure modes if the cache cluster throttles requests.
  • Broadcast state: For small, slowly changing lookup tables (such as postal code mappings or discount rules), broadcast the entire dataset to every parallel task slot. This avoids network lookups while keeping memory usage modest.

Temporal joins for point-in-time correctness

Standard lookup joins evaluate the customer record present in memory at processing time. If an order placed yesterday is delayed and arrives today, a simple lookup attaches today updated customer address instead of yesterday active address.

Temporal joins resolve this discrepancy by tracking the validity intervals of customer versions ([valid_from, valid_to)). The engine matches the order timestamp against the corresponding historical customer snapshot, producing deterministic, historically accurate enrichments.

PreviousNext