如何自定义KStream-KStream Join中的窗口存储实现?
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

