You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何自定义KStream-KStream Join中的窗口存储实现?

Customizing KStream-KStream Join Window Store to Retain Only Latest Records

Great question! Let's walk through how to adjust the window store used for your KStream-KStream join so that only the most recent record per key is kept in the window—this will ensure each side of the join only associates with the latest entry from the other stream.

The core issue here is that the default RocksDBWindowStore for joins has retainDuplicates=true, which preserves all records for a key within the window. To change this, you don't need to fully rewrite the store implementation; instead, you can create a custom store supplier with the desired behavior and attach it to your join. Here's how:

1. Create a Custom Window Store Supplier with retainDuplicates=false

Use Kafka Streams' built-in Stores utility to create a persistent (or in-memory) window store where older records for the same key are overwritten by newer ones. The critical parameter here is retainDuplicates=false:

// Define your window dimensions (match your large window requirement)
Duration windowSize = Duration.ofHours(24);
Duration advanceInterval = Duration.ofMinutes(5); // Adjust based on your needs

// Build the custom store supplier
WindowBytesStoreSupplier customJoinStore = Stores.persistentWindowStore(
    "latest-record-join-store", // Unique name for the store
    windowSize,
    advanceInterval,
    false, // Set to false to retain only the latest record per key
    new RocksDBConfigSetter() {
        // Optional: Add custom RocksDB configurations if needed
        @Override
        public void setConfig(String storeName, Options options, Map<String, Object> configs) {
            // Example: Enable Snappy compression for better storage efficiency
            options.setCompressionType(CompressionType.SNAPPY_COMPRESSION);
        }
    }
);

If you prefer an in-memory store (no persistence), replace persistentWindowStore with inMemoryWindowStore—it supports the same retainDuplicates parameter.

2. Attach the Custom Store to Your Join Operation

When configuring your KStream-KStream join, use the withStoreSupplier method to replace the default store with your custom one. This ensures both the left and right streams use the store that retains only latest records:

// Assume you have your two streams and serdes defined
KStream<String, LeftEvent> leftStream = ...;
KStream<String, RightEvent> rightStream = ...;
Serde<String> keySerde = Serdes.String();
Serde<LeftEvent> leftValueSerde = ...;
Serde<RightEvent> rightValueSerde = ...;

// Perform the join with the custom store
KStream<String, CombinedEvent> joinedStream = leftStream.join(
    rightStream,
    (leftEvent, rightEvent) -> new CombinedEvent(leftEvent, rightEvent), // Logic to combine values
    JoinWindows.of(windowSize), // Match your large window size
    Joined.with(keySerde, leftValueSerde, rightValueSerde)
          .withStoreSupplier(customJoinStore) // Attach your custom store here
);

3. Verify the Behavior

With retainDuplicates=false, the window store will automatically overwrite any existing record for the same key when a new one arrives. When the join processor fetches records via window.fetch(), it will only retrieve the latest entry for each key—exactly what you need to ensure joins only use the most recent events from both streams.

Key Things to Keep in Mind:

  • Retention Period: Ensure the store's retention time (derived from the window size by default) is at least as long as your window to avoid losing data before the window closes.
  • Business Logic Alignment: Confirm that discarding older records for the same key doesn't conflict with other parts of your pipeline—since you explicitly want only the latest records, this should be safe.
  • Store Uniqueness: Give your store a unique name to avoid conflicts with other stores in your Kafka Streams application.

内容的提问来源于stack exchange,提问作者Edmondo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 09:51:53