Stream-Stream Joins in Apache Flink and Kafka Streams

By Marcus Chen • • 9 min read

Joining two unbounded data streams in real time is one of the most demanding operations in stream processing. Unlike batch joins, where both datasets are finite and fully available, stream-stream joins must contend with events arriving out of order, variable latency, and the fundamental challenge of deciding when enough data has arrived to produce a result. Getting this wrong means either missing valid matches or accumulating state until your cluster runs out of memory.

This guide examines how Apache Flink and Kafka Streams handle stream-stream joins, the trade-offs each framework makes, and the practical patterns you need for production deployments. If you are building pipelines that feed into systems like ClickHouse materialized views or Snowflake dynamic tables, understanding these join semantics is essential.

Why Stream-Stream Joins Are Hard

In a traditional relational database, a join operates on two complete tables. The optimizer can scan both sides, build hash tables, and produce a deterministic result. Stream-stream joins have none of these luxuries. Events arrive continuously with no defined endpoint. An order event might arrive seconds, minutes, or even hours before the corresponding payment event that you want to join it with.

Three fundamental challenges define stream-stream joins:

  • Temporal alignment. Events from different streams arrive at different rates and with different delays. The join must decide how long to wait for a matching event before giving up.
  • State management. Both sides of the join must buffer events while waiting for potential matches. Without bounds on this buffering, state grows without limit.
  • Correctness under disorder. Late-arriving events can invalidate previously emitted results or create matches that should have been produced earlier.

Both Flink and Kafka Streams solve these problems, but with different abstractions and guarantees. The choice between them depends on your latency requirements, the complexity of your join conditions, and the operational model your team can support.

Join Types and Semantics

Before diving into framework-specific implementations, it helps to clarify the join types available in stream processing. Not every type that exists in SQL translates cleanly to the streaming world.

Join TypeFlink SupportKafka Streams SupportNotes
Inner JoinYesYesEmits when both sides match within the window
Left Outer JoinYesYesEmits unmatched left events when the window closes
Full Outer JoinYesYesEmits unmatched events from both sides
Interval JoinYesNo (use windowed)Per-event time boundaries rather than fixed windows
Temporal JoinYes (SQL API)KTable lookupJoin a stream against a versioned table

The distinction between windowed joins and interval joins is particularly important. A windowed join groups events into discrete time buckets, and events in the same bucket get joined. An interval join defines a time range relative to each individual event, offering finer control over match boundaries.

Windowed Joins in Kafka Streams

Kafka Streams provides a fluent API for stream-stream joins through the KStream.join() method. The framework handles state management internally using RocksDB-backed state stores, and the join window determines how long events are retained for matching.

Here is a concrete example joining an order stream with a payment stream, where payments are expected within 30 minutes of the order:

// Define the join window: payments within 30 minutes after the order
JoinWindows joinWindow = JoinWindows
    .ofTimeDifferenceWithNoGrace(Duration.ofMinutes(30));

// Perform the stream-stream join
KStream<String, EnrichedOrder> enrichedOrders = orderStream.join(
    paymentStream,
    (order, payment) -> new EnrichedOrder(
        order.getOrderId(),
        order.getCustomerId(),
        order.getAmount(),
        payment.getTransactionId(),
        payment.getPaymentMethod()
    ),
    joinWindow,
    StreamJoined.with(
        Serdes.String(),
        orderSerde,
        paymentSerde
    )
);

enrichedOrders.to("enriched-orders", Produced.with(
    Serdes.String(), enrichedOrderSerde
));

Several details matter in production. The ofTimeDifferenceWithNoGrace method means late-arriving events outside the window are silently dropped. If your use case requires handling late data, use ofTimeDifferenceAndGrace instead and specify a grace period:

JoinWindows joinWindow = JoinWindows
    .ofTimeDifferenceAndGrace(
        Duration.ofMinutes(30),  // join window
        Duration.ofMinutes(5)    // grace period for late arrivals
    );

State Store Considerations

Each side of the join maintains a windowed state store. The default RocksDB backend writes to local disk, which means your Kafka Streams instances need adequate disk space and IOPS. For a join window of 30 minutes with 10,000 events per second on each side, you are looking at roughly 18 million records in each state store at steady state.

Monitor the state store size through JMX metrics like rocksdb-state-id.estimate-num-keys and rocksdb-state-id.total-sst-files-size. Unexpected growth usually indicates that your join key has high cardinality with low match rates, meaning events accumulate without being joined and only get purged when the window expires.

Interval Joins in Apache Flink

Flink's interval join is one of its most powerful features for stream-stream joining. Unlike windowed joins that partition time into fixed buckets, interval joins define a time range relative to each event. This eliminates the boundary problem where two related events fall into adjacent windows and never get matched.

Here is the same order-payment join implemented with Flink's DataStream API:

DataStream<EnrichedOrder> enriched = orderStream
    .keyBy(order -> order.getOrderId())
    .intervalJoin(
        paymentStream.keyBy(payment -> payment.getOrderId())
    )
    .between(Time.minutes(-5), Time.minutes(30))
    .process(new ProcessJoinFunction<Order, Payment, EnrichedOrder>() {
        @Override
        public void processElement(
                Order order,
                Payment payment,
                Context ctx,
                Collector<EnrichedOrder> out) {
            out.collect(new EnrichedOrder(
                order.getOrderId(),
                order.getCustomerId(),
                order.getAmount(),
                payment.getTransactionId(),
                payment.getPaymentMethod()
            ));
        }
    });

The .between(Time.minutes(-5), Time.minutes(30)) clause means: for each order event at time t, match it with any payment event whose timestamp falls in the range [t - 5 minutes, t + 30 minutes]. The negative lower bound accounts for scenarios where a payment event might arrive slightly before the order event due to clock skew or processing delays.

Watermark Configuration

Watermarks are the mechanism Flink uses to track event-time progress. For interval joins, watermarks determine when buffered events can be safely discarded. A watermark of time W means: no events with timestamp less than W will arrive in the future.

DataStream<Order> orderStream = env
    .fromSource(kafkaSource, WatermarkStrategy
        .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(10))
        .withTimestampAssigner((order, ts) -> order.getEventTime())
        .withIdleness(Duration.ofMinutes(1)),
        "orders"
    );

The forBoundedOutOfOrderness strategy with a 10-second bound means watermarks lag 10 seconds behind the maximum observed timestamp. This gives late events a 10-second grace period. The withIdleness call prevents idle partitions from holding back the global watermark, which is a common issue when some Kafka partitions have low traffic.

Flink SQL for Stream Joins

For teams that prefer a declarative approach, Flink SQL offers interval joins with familiar syntax. This is particularly useful when your data engineers are more comfortable with SQL than Java:

SELECT
    o.order_id,
    o.customer_id,
    o.amount,
    p.transaction_id,
    p.payment_method,
    o.order_time
FROM orders o
JOIN payments p
    ON o.order_id = p.order_id
    AND p.payment_time BETWEEN o.order_time - INTERVAL '5' MINUTE
                          AND o.order_time + INTERVAL '30' MINUTE;

Flink's SQL planner translates this into the same interval join operator used by the DataStream API. The execution plan, state management, and watermark handling are identical. You can inspect the physical plan using EXPLAIN to verify that the optimizer chose the expected join strategy.

Temporal joins in Flink SQL offer another pattern worth knowing. When joining a stream against a slowly-changing dimension (like a currency exchange rate table), you can use a temporal join to look up the rate that was valid at the time of each transaction. This avoids the need for a separate lookup service and keeps all join logic within the streaming pipeline.

Production Considerations and State Management

Running stream-stream joins in production is fundamentally a state management problem. Here are the key areas that demand attention:

State Size and Checkpointing

Both frameworks maintain state proportional to the join window size multiplied by the event rate. For Flink, enable incremental checkpoints with the RocksDB state backend to avoid checkpoint duration scaling linearly with state size:

state.backend: rocksdb
state.backend.incremental: true
state.checkpoints.dir: s3://your-bucket/flink/checkpoints
execution.checkpointing.interval: 60000
execution.checkpointing.min-pause: 30000

For Kafka Streams, ensure that standby replicas are configured so that failover does not require a full state store rebuild from the changelog topic:

num.standby.replicas=1
state.dir=/mnt/fast-ssd/kafka-streams

Handling Late Data

Late data in stream-stream joins creates a dilemma. Allowing late arrivals to produce join results means downstream consumers see updates to previously-emitted windows, which complicates exactly-once semantics. Dropping late data means accepting some level of incompleteness.

In practice, most teams take a layered approach. The streaming join handles events within a reasonable window, say 30 minutes plus a 5-minute grace period. Events that arrive later are captured in a dead letter queue and reconciled through a batch job that runs daily. This combination gives you real-time results for the vast majority of events while still achieving eventual completeness. The reconciled results can then feed downstream systems like Delta Lake change data feeds for reliable propagation to analytics consumers.

Key Skew and Hot Keys

If certain join keys receive disproportionately more traffic, one parallel instance handles far more state and computation than others. In Flink, you can use a two-phase join strategy: first broadcast the smaller stream, then join locally without shuffle. In Kafka Streams, you can pre-partition the skewed key space across multiple sub-keys and aggregate after the join.

Choosing Between Flink and Kafka Streams

The decision between Flink and Kafka Streams for stream-stream joins is not purely technical. It involves operational, organizational, and architectural considerations.

ConsiderationApache FlinkKafka Streams
DeploymentStandalone cluster or YARN/K8sEmbedded in your application
Join sophisticationInterval, temporal, windowed, patternWindowed joins with configurable grace
Event-time handlingFirst-class watermark supportTimestamp extractor + grace periods
State backendRocksDB, HashMaps, customRocksDB or in-memory
Exactly-onceVia checkpoints + 2PC sinksVia Kafka transactions
Learning curveSteeper (distributed system)Gentler (library, not framework)

Choose Kafka Streams when your joins are straightforward windowed operations, your source and sink are both Kafka, and you value operational simplicity. Choose Flink when you need interval joins, complex event processing, or need to join streams from heterogeneous sources beyond Kafka.

Teams already running Flink for other workloads should lean toward consolidating stream joins onto the same platform rather than introducing Kafka Streams as an additional technology. Conversely, if your infrastructure is Kafka-centric and you do not want to operate a separate compute cluster, Kafka Streams keeps your stack simpler.

Testing Stream-Stream Joins

Testing joins in streaming pipelines requires a different approach than testing batch transformations. You need to verify behavior under specific timing conditions, not just data correctness.

For Kafka Streams, the TopologyTestDriver provides a deterministic testing harness that lets you control event timestamps precisely:

TopologyTestDriver testDriver = new TopologyTestDriver(
    topology, streamProperties
);

TestInputTopic<String, Order> orderInput = testDriver
    .createInputTopic("orders", stringSerde, orderSerde);
TestInputTopic<String, Payment> paymentInput = testDriver
    .createInputTopic("payments", stringSerde, paymentSerde);
TestOutputTopic<String, EnrichedOrder> output = testDriver
    .createOutputTopic("enriched-orders", stringSerde, enrichedSerde);

// Simulate an order followed by a payment 10 minutes later
Instant orderTime = Instant.parse("2026-10-01T10:00:00Z");
orderInput.pipeInput("order-1", testOrder, orderTime);

Instant paymentTime = orderTime.plus(Duration.ofMinutes(10));
paymentInput.pipeInput("order-1", testPayment, paymentTime);

// Verify the join produced a result
assertThat(output.readValue().getTransactionId())
    .isEqualTo(testPayment.getTransactionId());

For Flink, the MiniClusterExtension provides a local Flink environment for integration tests. You can use the test harness to inject watermarks at specific points and verify that late events are handled correctly.

In both cases, write tests that specifically cover edge cases: events arriving exactly at the window boundary, events that arrive out of order, and events that arrive after the grace period has expired. These are the scenarios that cause production incidents, and they are the ones that unit tests on business logic alone will never catch.