This is an automated email from the ASF dual-hosted git repository.

je-ik pushed a commit to branch feat/18479-kafka-streams-runner-skeleton
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to 
refs/heads/feat/18479-kafka-streams-runner-skeleton by this push:
     new 47614a943a7 [GSoC 2026] Kafka Streams runner: windowed GroupByKey via 
ReduceFnRunner (#39494)
47614a943a7 is described below

commit 47614a943a7046b77056e9fd087b0c0482f24ac6
Author: M Junaid Shaukat <[email protected]>
AuthorDate: Tue Jul 28 13:12:59 2026 +0500

    [GSoC 2026] Kafka Streams runner: windowed GroupByKey via ReduceFnRunner 
(#39494)
    
    * [GSoC 2026] Kafka Streams runner: windowed GroupByKey via ReduceFnRunner
    
    Replaces the global-window-only GroupByKey with a windowed one that drives
    Beam's ReduceFnRunner on the runner side, the way the Flink and Spark
    portable runners do, backed by Kafka Streams state and timers.
    
    WindowedGroupByKeyProcessor builds a ReduceFnRunner per key (like
    GroupAlsoByWindowViaWindowSetNewDoFn) over two new backends:
    KafkaStreamsStateInternals, which stores each Beam state cell as one entry
    in a KeyValueStore under a composite key of key + namespace + tag
    (modelled on SparkStateInternals), and KafkaStreamsTimerInternals, which
    persists timers keyed by identity and is fired by the processor scanning
    for due event-time timers on each input-watermark advance.
    GroupByKeyTranslator hydrates the input windowing strategy from the
    pipeline proto and wires the state and timer stores. Windowing, the
    default trigger, panes, allowed lateness and timestamp combiners all come
    from ReduceFnRunner.
---
 runners/kafka-streams/build.gradle                 |  11 +-
 .../streams/translation/GroupByKeyProcessor.java   | 224 ----------
 .../streams/translation/GroupByKeyTranslator.java  |  78 +++-
 .../translation/KafkaStreamsStateInternals.java    | 463 +++++++++++++++++++++
 .../translation/KafkaStreamsTimerInternals.java    | 265 ++++++++++++
 .../kafka/streams/translation/StoreKeys.java       | 105 +++++
 .../translation/WindowedGroupByKeyProcessor.java   | 323 ++++++++++++++
 .../kafka/streams/KafkaStreamsTestRunner.java      | 101 +----
 .../translation/FixedWindowGroupByKeyTest.java     | 106 +++++
 .../KafkaStreamsTimerInternalsTest.java            | 208 +++++++++
 10 files changed, 1552 insertions(+), 332 deletions(-)

diff --git a/runners/kafka-streams/build.gradle 
b/runners/kafka-streams/build.gradle
index 2d0209eafbf..bdf7e3be058 100644
--- a/runners/kafka-streams/build.gradle
+++ b/runners/kafka-streams/build.gradle
@@ -87,11 +87,11 @@ dependencies {
 
 // Known-failing @ValidatesRunner tests, excluded until the feature they need 
lands.
 def sickbayTests = [
-  // Non-global windowing (FixedWindows, merging windows, timestamp combiners) 
is not supported
-  // yet; these apply a window and assert on window-derived output, hitting a 
GlobalWindow cast.
-  'org.apache.beam.sdk.transforms.GroupByKeyTest$WindowTests',
-  
'org.apache.beam.sdk.transforms.GroupByKeyTest$BasicTests.testTimestampCombinerLatest',
-  
'org.apache.beam.sdk.transforms.GroupByKeyTest$BasicTests.testTimestampCombinerEarliest',
+  // Merging (session) windows are not supported yet: ReduceFnRunner drives 
them through a merging
+  // window set that moves per-window state as windows merge, which this first 
windowing pass does
+  // not implement. Non-merging windows (fixed, sliding), the default trigger 
and timestamp
+  // combiners do work. Lands with the follow-up windowing PR.
+  
'org.apache.beam.sdk.transforms.GroupByKeyTest$WindowTests.testGroupByKeyMergingWindows',
   // A DoFn whose @StartBundle throws never gets to report its error: 
SdkHarnessClient.newBundle
   // sends the ProcessBundleRequest and then blocks in 
GrpcDataService.createOutboundAggregator
   // waiting for the SDK harness to open its data stream, which a bundle that 
failed during setup
@@ -100,6 +100,7 @@ def sickbayTests = [
   // anything specific to this runner; the Flink runner sickbays all of 
LifecycleTests and the
   // Prism runner sickbays each of its three error tests. The @ProcessElement 
and @FinishBundle
   // variants do pass here, because by then the data stream is established.
+  // Tracked by https://github.com/apache/beam/issues/39452.
   
'org.apache.beam.sdk.transforms.ParDoTest$LifecycleTests.testParDoWithErrorInStartBatch',
 ]
 
diff --git 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyProcessor.java
 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyProcessor.java
deleted file mode 100644
index 3e82935b807..00000000000
--- 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyProcessor.java
+++ /dev/null
@@ -1,224 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.apache.beam.runners.kafka.streams.translation;
-
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Set;
-import org.apache.beam.sdk.coders.Coder;
-import org.apache.beam.sdk.coders.CoderException;
-import org.apache.beam.sdk.coders.IterableCoder;
-import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
-import org.apache.beam.sdk.transforms.windowing.GlobalWindow;
-import org.apache.beam.sdk.util.CoderUtils;
-import org.apache.beam.sdk.values.KV;
-import org.apache.beam.sdk.values.WindowedValue;
-import org.apache.beam.sdk.values.WindowedValues;
-import org.apache.kafka.streams.processor.api.Processor;
-import org.apache.kafka.streams.processor.api.ProcessorContext;
-import org.apache.kafka.streams.processor.api.Record;
-import org.apache.kafka.streams.state.KeyValueIterator;
-import org.apache.kafka.streams.state.KeyValueStore;
-import org.checkerframework.checker.nullness.qual.Nullable;
-import org.joda.time.Instant;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-/**
- * Executes a {@code GroupByKey} (GlobalWindow, default trigger, no allowed 
lateness).
- *
- * <p>Records arrive on the repartition topic keyed by the encoded Beam key, 
so every value of a key
- * is co-located here. Each value is appended to a per-key buffer in a Kafka 
Streams state store.
- * Watermark reports are fed to a {@link WatermarkAggregator}; when the input 
watermark reaches
- * {@link BoundedWindow#TIMESTAMP_MAX_VALUE} (the end of the global window) 
every buffered key is
- * emitted once as {@code KV<K, Iterable<V>>} and the buffer cleared, then the 
watermark is
- * forwarded downstream.
- *
- * <p>Buffering whole value lists and re-encoding on each append is O(n^2) per 
key; fine for this
- * first GroupByKey, and replaced when this moves to runner-core {@code 
GroupAlsoByWindow}.
- */
-class GroupByKeyProcessor
-    implements Processor<byte[], KStreamsPayload<?>, byte[], 
KStreamsPayload<?>> {
-
-  private static final Logger LOG = 
LoggerFactory.getLogger(GroupByKeyProcessor.class);
-
-  private final String stateStoreName;
-  // This transform's own id, stamped on every watermark it forwards 
downstream.
-  private final String transformId;
-  private final Coder<Object> keyCoder;
-  private final IterableCoder<@Nullable Object> bufferCoder;
-
-  // Aggregates the input watermark from the upstream transform's reports, 
which arrive through the
-  // repartition topic with the upstream producer's transform id intact (the 
shuffle forwards
-  // watermark payloads unchanged).
-  private final WatermarkAggregator watermarkAggregator;
-  private Instant lastForwardedWatermark = BoundedWindow.TIMESTAMP_MIN_VALUE;
-  // The global window fires exactly once, when the watermark first reaches 
its end. Later watermark
-  // reports (e.g. the same terminal watermark broadcast across repartition 
partitions) must not
-  // re-fire. This flag is in-memory only; restart correctness comes from the 
state store plus
-  // exactly-once-v2: the buffered values and consumer offsets are committed 
atomically, and the
-  // store is empty once a key has fired, so a restart cannot double-emit. 
Persisting watermark
-  // holds is part of the separate WatermarkManager persistence work, not this 
initial GroupByKey.
-  private boolean fired = false;
-
-  private @Nullable ProcessorContext<byte[], KStreamsPayload<?>> context;
-  private @Nullable KeyValueStore<byte[], byte[]> store;
-
-  /**
-   * @param transformId this transform's own id, stamped on the watermarks it 
emits
-   * @param upstreamTransformIds the transform ids feeding this GroupByKey 
(known from the pipeline
-   *     graph), whose reports the {@link WatermarkAggregator} waits for
-   */
-  GroupByKeyProcessor(
-      String stateStoreName,
-      String transformId,
-      Set<String> upstreamTransformIds,
-      Coder<Object> keyCoder,
-      Coder<@Nullable Object> valueCoder) {
-    this.stateStoreName = stateStoreName;
-    this.transformId = transformId;
-    this.watermarkAggregator = new WatermarkAggregator(upstreamTransformIds);
-    this.keyCoder = keyCoder;
-    this.bufferCoder = IterableCoder.of(valueCoder);
-  }
-
-  @Override
-  public void init(ProcessorContext<byte[], KStreamsPayload<?>> context) {
-    this.context = context;
-    this.store = context.getStateStore(stateStoreName);
-  }
-
-  @Override
-  public void process(Record<byte[], KStreamsPayload<?>> record) {
-    KStreamsPayload<?> payload = record.value();
-    if (payload == null) {
-      // The repartition topic can be written to from outside the runner (or 
carry a tombstone),
-      // so recover from the obvious error instead of crashing the task: warn 
and drop.
-      LOG.warn(
-          "GroupByKey {} dropping record with null payload (external write or 
tombstone)",
-          transformId);
-      return;
-    }
-    if (payload.isData()) {
-      byte[] encodedKey = record.key();
-      Object element = payload.getData().getValue();
-      if (encodedKey == null || element == null) {
-        throw new IllegalStateException("GroupByKey data record is missing its 
key or value");
-      }
-      appendValue(encodedKey, element);
-      return;
-    }
-    watermarkAggregator.observe(payload.asWatermark());
-    Instant advanced = watermarkAggregator.advance();
-    if (!fired && !advanced.isBefore(BoundedWindow.TIMESTAMP_MAX_VALUE)) {
-      fireAll(record);
-      fired = true;
-    }
-    if (advanced.isAfter(lastForwardedWatermark)) {
-      lastForwardedWatermark = advanced;
-      forwardWatermark(record, advanced.getMillis());
-    }
-  }
-
-  private void appendValue(byte[] encodedKey, Object kvObject) {
-    KV<?, ?> kv = (KV<?, ?>) kvObject;
-    KeyValueStore<byte[], byte[]> kvStore = checkInitialized(store);
-    byte[] existing = kvStore.get(encodedKey);
-    List<@Nullable Object> values = existing == null ? new ArrayList<>() : 
decodeBuffer(existing);
-    values.add(kv.getValue());
-    kvStore.put(encodedKey, encodeBuffer(values));
-  }
-
-  private void fireAll(Record<byte[], KStreamsPayload<?>> trigger) {
-    // NOTE: this emits every buffered key in a single watermark turn. For a 
very large key space
-    // that risks memory pressure and exceeding the poll / transaction 
timeout. Acceptable for this
-    // initial GlobalWindow GroupByKey (fire once at end of input); 
incremental, timer-driven output
-    // via runner-core GroupAlsoByWindow lands with the windowing/timers work.
-    ProcessorContext<byte[], KStreamsPayload<?>> ctx = 
checkInitialized(context);
-    KeyValueStore<byte[], byte[]> kvStore = checkInitialized(store);
-    List<byte[]> firedKeys = new ArrayList<>();
-    try (KeyValueIterator<byte[], byte[]> it = kvStore.all()) {
-      while (it.hasNext()) {
-        org.apache.kafka.streams.KeyValue<byte[], byte[]> entry = it.next();
-        Object key = decodeKey(entry.key);
-        List<@Nullable Object> values = decodeBuffer(entry.value);
-        // The pane fires at the end of the global window, so the grouped 
element carries the
-        // window's max timestamp (END_OF_GLOBAL_WINDOW). Emitting at 
TIMESTAMP_MIN_VALUE (the
-        // default of valueInGlobalWindow) would make the output appear 
arbitrarily late and be
-        // dropped downstream once the watermark has advanced.
-        WindowedValue<KV<Object, Iterable<@Nullable Object>>> output =
-            WindowedValues.timestampedValueInGlobalWindow(
-                KV.of(key, (Iterable<@Nullable Object>) values),
-                GlobalWindow.INSTANCE.maxTimestamp());
-        ctx.forward(
-            new Record<byte[], KStreamsPayload<?>>(
-                entry.key, KStreamsPayload.data(output), trigger.timestamp()));
-        firedKeys.add(entry.key);
-      }
-    }
-    for (byte[] key : firedKeys) {
-      kvStore.delete(key);
-    }
-  }
-
-  private void forwardWatermark(Record<byte[], KStreamsPayload<?>> trigger, 
long watermarkMillis) {
-    ProcessorContext<byte[], KStreamsPayload<?>> ctx = 
checkInitialized(context);
-    // Stamped with this transform's own id; GroupByKey is a single instance 
for now, so the report
-    // is for its only partition (0 of 1).
-    ctx.forward(
-        new Record<byte[], KStreamsPayload<?>>(
-            trigger.key(),
-            KStreamsPayload.watermark(watermarkMillis, transformId, 0, 1),
-            trigger.timestamp()));
-  }
-
-  private byte[] encodeBuffer(List<@Nullable Object> values) {
-    try {
-      return CoderUtils.encodeToByteArray(bufferCoder, values);
-    } catch (CoderException e) {
-      throw new RuntimeException("Failed to encode GroupByKey value buffer", 
e);
-    }
-  }
-
-  private List<@Nullable Object> decodeBuffer(byte[] bytes) {
-    try {
-      List<@Nullable Object> values = new ArrayList<>();
-      for (@Nullable Object value : 
CoderUtils.decodeFromByteArray(bufferCoder, bytes)) {
-        values.add(value);
-      }
-      return values;
-    } catch (CoderException e) {
-      throw new RuntimeException("Failed to decode GroupByKey value buffer", 
e);
-    }
-  }
-
-  private Object decodeKey(byte[] bytes) {
-    try {
-      return CoderUtils.decodeFromByteArray(keyCoder, bytes);
-    } catch (CoderException e) {
-      throw new RuntimeException("Failed to decode GroupByKey key", e);
-    }
-  }
-
-  private static <T> T checkInitialized(@Nullable T value) {
-    if (value == null) {
-      throw new IllegalStateException("GroupByKeyProcessor used before 
init()");
-    }
-    return value;
-  }
-}
diff --git 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
index 9e23dbb5cfb..c5327e28e06 100644
--- 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
+++ 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
@@ -22,8 +22,12 @@ import static 
org.apache.beam.runners.fnexecution.translation.PipelineTranslator
 import org.apache.beam.model.pipeline.v1.RunnerApi;
 import org.apache.beam.sdk.coders.Coder;
 import org.apache.beam.sdk.coders.KvCoder;
+import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.util.construction.RehydratedComponents;
+import org.apache.beam.sdk.util.construction.WindowingStrategyTranslation;
 import org.apache.beam.sdk.values.KV;
 import org.apache.beam.sdk.values.WindowedValues;
+import org.apache.beam.sdk.values.WindowingStrategy;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
 import org.apache.kafka.common.serialization.Serdes;
@@ -35,10 +39,11 @@ import org.checkerframework.checker.nullness.qual.Nullable;
  * Translates the {@code beam:transform:group_by_key:v1} URN — the runner's 
first stateful,
  * shuffle-bearing transform.
  *
- * <p>This is the simplest GroupByKey: GlobalWindow, default trigger, no 
allowed lateness (per the
- * plan agreed with the mentor). Each key's values are buffered in a Kafka 
Streams state store and
- * emitted once as {@code KV<K, Iterable<V>>} when the watermark reaches {@link
- * org.apache.beam.sdk.transforms.windowing.BoundedWindow#TIMESTAMP_MAX_VALUE}.
+ * <p>Windowing and triggering are executed by Beam's {@link
+ * org.apache.beam.runners.core.ReduceFnRunner} inside {@link 
WindowedGroupByKeyProcessor}, the same
+ * way the Flink and Spark portable runners do it — so fixed/sliding windows, 
the default trigger,
+ * allowed lateness and timestamp combiners all work. The input PCollection's 
windowing strategy is
+ * hydrated from the pipeline proto and handed to the processor.
  *
  * <p>Topology added (the Beam key becomes the Kafka record key so Kafka 
Streams shuffles by it):
  *
@@ -49,7 +54,8 @@ import org.checkerframework.checker.nullness.qual.Nullable;
  *       via {@link KStreamsPayloadSerde} and a {@link 
GroupByKeyBroadcastPartitioner} that hashes
  *       data by key and fans watermark reports out to every partition;
  *   <li>a {@link Topology#addSource source} reading the repartition topic 
back;
- *   <li>the {@link GroupByKeyProcessor} plus a persistent state store, wired 
to the source.
+ *   <li>the {@link WindowedGroupByKeyProcessor} plus persistent state and 
timer stores, wired to
+ *       the source.
  * </ul>
  *
  * <p>The repartition topic is expected to exist on the broker before the job 
starts (same
@@ -62,6 +68,9 @@ class GroupByKeyTranslator implements PTransformTranslator {
   static final String SINK_SUFFIX = "-repartition-sink";
   static final String SOURCE_SUFFIX = "-repartition-source";
   static final String STATE_STORE_SUFFIX = "-state";
+  static final String HOLDS_INDEX_STORE_SUFFIX = "-holds-index";
+  static final String TIMER_STORE_SUFFIX = "-timers";
+  static final String TIMER_INDEX_STORE_SUFFIX = "-timers-index";
   static final String REPARTITION_TOPIC_PREFIX = "__beam_gbk_";
 
   @Override
@@ -82,12 +91,18 @@ class GroupByKeyTranslator implements PTransformTranslator {
     Coder<@Nullable Object> valueCoder =
         (Coder<@Nullable Object>) (Coder<?>) kvCoder.getValueCoder();
 
+    WindowingStrategy<?, BoundedWindow> windowingStrategy =
+        hydrateWindowingStrategy(pipeline, inputPCollectionId);
+
     String parentProcessor = 
context.getProcessorNameForPCollection(inputPCollectionId);
 
     String shuffleName = transformId + SHUFFLE_SUFFIX;
     String sinkName = transformId + SINK_SUFFIX;
     String sourceName = transformId + SOURCE_SUFFIX;
     String stateStoreName = transformId + STATE_STORE_SUFFIX;
+    String holdsIndexStoreName = transformId + HOLDS_INDEX_STORE_SUFFIX;
+    String timerStoreName = transformId + TIMER_STORE_SUFFIX;
+    String timerIndexStoreName = transformId + TIMER_INDEX_STORE_SUFFIX;
     String repartitionTopic = repartitionTopic(transformId);
 
     KStreamsPayloadSerde<KV<Object, Object>> payloadSerde = new 
KStreamsPayloadSerde<>(inputCoder);
@@ -111,27 +126,70 @@ class GroupByKeyTranslator implements 
PTransformTranslator {
         payloadSerde.deserializer(),
         repartitionTopic);
 
-    // Buffer values per key and fire KV<K, Iterable<V>> at the terminal 
watermark. Watermark
-    // reports cross the repartition topic unchanged, so they still carry the 
id of the transform
-    // that produced this GroupByKey's input — the parent the shuffle is 
attached to.
+    // Group by key and window through Beam's ReduceFnRunner, backed by the 
state and timer stores.
+    // Watermark reports cross the repartition topic unchanged, so they still 
carry the id of the
+    // transform that produced this GroupByKey's input — the parent the 
shuffle is attached to.
     topology.addProcessor(
         transformId,
         () ->
-            new GroupByKeyProcessor(
+            new WindowedGroupByKeyProcessor<Object, @Nullable Object, 
BoundedWindow>(
                 stateStoreName,
+                holdsIndexStoreName,
+                timerStoreName,
+                timerIndexStoreName,
                 transformId,
                 ImmutableSet.of(parentProcessor),
                 keyCoder,
-                valueCoder),
+                valueCoder,
+                windowingStrategy,
+                context.getPipelineOptions()),
         sourceName);
     topology.addStateStore(
         Stores.keyValueStoreBuilder(
             Stores.persistentKeyValueStore(stateStoreName), 
Serdes.ByteArray(), Serdes.ByteArray()),
         transformId);
+    topology.addStateStore(
+        Stores.keyValueStoreBuilder(
+            Stores.persistentKeyValueStore(timerStoreName), 
Serdes.ByteArray(), Serdes.ByteArray()),
+        transformId);
+    // Indexes ordered by timestamp, so due timers and the minimum watermark 
hold are range scans
+    // rather than scans of every timer or every held window.
+    topology.addStateStore(
+        Stores.keyValueStoreBuilder(
+            Stores.persistentKeyValueStore(timerIndexStoreName),
+            Serdes.ByteArray(),
+            Serdes.ByteArray()),
+        transformId);
+    topology.addStateStore(
+        Stores.keyValueStoreBuilder(
+            Stores.persistentKeyValueStore(holdsIndexStoreName),
+            Serdes.ByteArray(),
+            Serdes.ByteArray()),
+        transformId);
 
     context.registerPCollectionProducer(outputPCollectionId, transformId);
   }
 
+  /** Hydrates the input PCollection's windowing strategy from the pipeline 
proto. */
+  private static WindowingStrategy<?, BoundedWindow> hydrateWindowingStrategy(
+      RunnerApi.Pipeline pipeline, String inputPCollectionId) {
+    RunnerApi.Components components = pipeline.getComponents();
+    String windowingStrategyId =
+        
components.getPcollectionsOrThrow(inputPCollectionId).getWindowingStrategyId();
+    try {
+      @SuppressWarnings("unchecked")
+      WindowingStrategy<?, BoundedWindow> strategy =
+          (WindowingStrategy<?, BoundedWindow>)
+              WindowingStrategyTranslation.fromProto(
+                  
components.getWindowingStrategiesOrThrow(windowingStrategyId),
+                  RehydratedComponents.forComponents(components));
+      return strategy;
+    } catch (Exception e) {
+      throw new IllegalStateException(
+          "Failed to hydrate GroupByKey windowing strategy " + 
windowingStrategyId, e);
+    }
+  }
+
   /** The internal repartition topic name for a GroupByKey transform. */
   static String repartitionTopic(String transformId) {
     return REPARTITION_TOPIC_PREFIX + 
transformId.replaceAll("[^a-zA-Z0-9._-]", "_");
diff --git 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsStateInternals.java
 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsStateInternals.java
new file mode 100644
index 00000000000..c17d7c45ad4
--- /dev/null
+++ 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsStateInternals.java
@@ -0,0 +1,463 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.translation;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import org.apache.beam.runners.core.StateInternals;
+import org.apache.beam.runners.core.StateNamespace;
+import org.apache.beam.runners.core.StateTag;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.coders.CoderException;
+import org.apache.beam.sdk.coders.InstantCoder;
+import org.apache.beam.sdk.coders.ListCoder;
+import org.apache.beam.sdk.state.BagState;
+import org.apache.beam.sdk.state.CombiningState;
+import org.apache.beam.sdk.state.MapState;
+import org.apache.beam.sdk.state.MultimapState;
+import org.apache.beam.sdk.state.OrderedListState;
+import org.apache.beam.sdk.state.ReadableState;
+import org.apache.beam.sdk.state.SetState;
+import org.apache.beam.sdk.state.State;
+import org.apache.beam.sdk.state.StateBinder;
+import org.apache.beam.sdk.state.StateContext;
+import org.apache.beam.sdk.state.StateSpec;
+import org.apache.beam.sdk.state.ValueState;
+import org.apache.beam.sdk.state.WatermarkHoldState;
+import org.apache.beam.sdk.transforms.Combine.CombineFn;
+import org.apache.beam.sdk.transforms.CombineWithContext;
+import org.apache.beam.sdk.transforms.windowing.TimestampCombiner;
+import org.apache.beam.sdk.util.CoderUtils;
+import org.apache.beam.sdk.util.CombineFnUtil;
+import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.KeyValueStore;
+import org.checkerframework.checker.nullness.qual.NonNull;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.joda.time.Instant;
+
+/**
+ * A {@link StateInternals} for one key, backed by a Kafka Streams {@link 
KeyValueStore}.
+ *
+ * <p>Beam addresses a state cell by {@code (key, StateNamespace, StateTag)}; 
a windowed pipeline
+ * puts each window's state in its own namespace. Every cell is stored as one 
entry in the shared
+ * per-transform store under a composite byte key {@code len(key)|key | 
len(ns)|ns | len(tag)|tag},
+ * so all cells for one Beam key share a prefix and a whole key's state can be 
range-scanned. The
+ * value is the cell's contents encoded with its Beam {@link Coder}. Writing 
straight to the store
+ * (rather than buffering and flushing) keeps this restart-safe for free: the 
store is changelogged
+ * and, under exactly-once, its writes commit atomically with the input 
offsets.
+ *
+ * <p>Modeled on the Spark runner's {@code SparkStateInternals}; the 
difference is that each cell
+ * reads and writes its own store entry instead of an in-memory table, so 
there is no separate
+ * persist step.
+ */
+class KafkaStreamsStateInternals<K> implements StateInternals {
+
+  /** The holds index is a set; only its keys carry information. */
+  private static final byte[] EMPTY_VALUE = new byte[0];
+
+  /**
+   * Reads the minimum watermark hold held by any key and window, or {@code 
null} if none is held.
+   * The index is ordered by hold time, so this is the first entry rather than 
a scan.
+   */
+  static @Nullable Instant minWatermarkHold(KeyValueStore<byte[], byte[]> 
holdsIndexStore) {
+    try (KeyValueIterator<byte[], byte[]> it = holdsIndexStore.all()) {
+      if (!it.hasNext()) {
+        return null;
+      }
+      return new Instant(StoreKeys.readTimestamp(it.next().key, 0));
+    }
+  }
+
+  private final @NonNull K key;
+  private final byte[] encodedKey;
+  private final KeyValueStore<byte[], byte[]> store;
+  private final KeyValueStore<byte[], byte[]> holdsIndexStore;
+
+  /**
+   * The last namespace a composite key was built for, and the {@code key | 
namespace} prefix that
+   * was built for it. One turn of the windowing runner touches several tags 
in the same namespace
+   * back to back (read the buffer, read the hold, write both), so caching the 
prefix removes most
+   * of the per-access encoding work.
+   */
+  private @Nullable StateNamespace cachedNamespace;
+
+  private byte @Nullable [] cachedPrefix;
+
+  KafkaStreamsStateInternals(
+      @NonNull K key,
+      byte[] encodedKey,
+      KeyValueStore<byte[], byte[]> store,
+      KeyValueStore<byte[], byte[]> holdsIndexStore) {
+    this.key = key;
+    this.encodedKey = encodedKey;
+    this.store = store;
+    this.holdsIndexStore = holdsIndexStore;
+  }
+
+  @Override
+  public Object getKey() {
+    return key;
+  }
+
+  @Override
+  public <T extends State> T state(
+      StateNamespace namespace, StateTag<T> address, StateContext<?> c) {
+    return address.getSpec().bind(address.getId(), new 
KafkaStreamsStateBinder(namespace, c));
+  }
+
+  /**
+   * The composite store key for one cell: {@code len|key len|namespace 
len|tagId}.
+   *
+   * <p>Built into one exactly-sized array, reusing the cached {@code key | 
namespace} prefix. This
+   * runs on every state access, so it avoids the repeated growth and final 
copy a stream would do.
+   */
+  private byte[] compositeKey(StateNamespace namespace, String id) {
+    byte[] prefix = prefixFor(namespace);
+    byte[] idBytes = id.getBytes(StandardCharsets.UTF_8);
+    byte[] compositeKey = new byte[prefix.length + 
StoreKeys.segmentLength(idBytes)];
+    System.arraycopy(prefix, 0, compositeKey, 0, prefix.length);
+    StoreKeys.writeSegment(compositeKey, prefix.length, idBytes);
+    return compositeKey;
+  }
+
+  /** The {@code key | namespace} prefix every cell in {@code namespace} 
starts with. */
+  private byte[] prefixFor(StateNamespace namespace) {
+    byte[] cached = cachedPrefix;
+    if (cached != null && namespace.equals(cachedNamespace)) {
+      return cached;
+    }
+    byte[] namespaceBytes = 
namespace.stringKey().getBytes(StandardCharsets.UTF_8);
+    byte[] prefix =
+        new byte[StoreKeys.segmentLength(encodedKey) + 
StoreKeys.segmentLength(namespaceBytes)];
+    int offset = StoreKeys.writeSegment(prefix, 0, encodedKey);
+    StoreKeys.writeSegment(prefix, offset, namespaceBytes);
+    cachedNamespace = namespace;
+    cachedPrefix = prefix;
+    return prefix;
+  }
+
+  private class KafkaStreamsStateBinder implements StateBinder {
+    private final StateNamespace namespace;
+    private final StateContext<?> stateContext;
+
+    private KafkaStreamsStateBinder(StateNamespace namespace, StateContext<?> 
stateContext) {
+      this.namespace = namespace;
+      this.stateContext = stateContext;
+    }
+
+    @Override
+    public <T> ValueState<T> bindValue(String id, StateSpec<ValueState<T>> 
spec, Coder<T> coder) {
+      return new KafkaStreamsValueState<>(namespace, id, coder);
+    }
+
+    @Override
+    public <T> BagState<T> bindBag(String id, StateSpec<BagState<T>> spec, 
Coder<T> elemCoder) {
+      return new KafkaStreamsBagState<>(namespace, id, elemCoder);
+    }
+
+    @Override
+    public <T> SetState<T> bindSet(String id, StateSpec<SetState<T>> spec, 
Coder<T> elemCoder) {
+      throw new UnsupportedOperationException(
+          SetState.class.getSimpleName() + " is not supported by the Kafka 
Streams runner yet");
+    }
+
+    @Override
+    public <KeyT, ValueT> MapState<KeyT, ValueT> bindMap(
+        String id,
+        StateSpec<MapState<KeyT, ValueT>> spec,
+        Coder<KeyT> mapKeyCoder,
+        Coder<ValueT> mapValueCoder) {
+      throw new UnsupportedOperationException(
+          MapState.class.getSimpleName() + " is not supported by the Kafka 
Streams runner yet");
+    }
+
+    @Override
+    public <KeyT, ValueT> MultimapState<KeyT, ValueT> bindMultimap(
+        String id,
+        StateSpec<MultimapState<KeyT, ValueT>> spec,
+        Coder<KeyT> keyCoder,
+        Coder<ValueT> valueCoder) {
+      throw new UnsupportedOperationException(
+          MultimapState.class.getSimpleName()
+              + " is not supported by the Kafka Streams runner yet");
+    }
+
+    @Override
+    public <T> OrderedListState<T> bindOrderedList(
+        String id, StateSpec<OrderedListState<T>> spec, Coder<T> elemCoder) {
+      throw new UnsupportedOperationException(
+          OrderedListState.class.getSimpleName()
+              + " is not supported by the Kafka Streams runner yet");
+    }
+
+    @Override
+    public <InputT, AccumT, OutputT> CombiningState<InputT, AccumT, OutputT> 
bindCombining(
+        String id,
+        StateSpec<CombiningState<InputT, AccumT, OutputT>> spec,
+        Coder<AccumT> accumCoder,
+        CombineFn<InputT, AccumT, OutputT> combineFn) {
+      return new KafkaStreamsCombiningState<>(namespace, id, accumCoder, 
combineFn);
+    }
+
+    @Override
+    public <InputT, AccumT, OutputT>
+        CombiningState<InputT, AccumT, OutputT> bindCombiningWithContext(
+            String id,
+            StateSpec<CombiningState<InputT, AccumT, OutputT>> spec,
+            Coder<AccumT> accumCoder,
+            CombineWithContext.CombineFnWithContext<InputT, AccumT, OutputT> 
combineFn) {
+      return new KafkaStreamsCombiningState<>(
+          namespace, id, accumCoder, CombineFnUtil.bindContext(combineFn, 
stateContext));
+    }
+
+    @Override
+    public WatermarkHoldState bindWatermark(
+        String id, StateSpec<WatermarkHoldState> spec, TimestampCombiner 
timestampCombiner) {
+      return new KafkaStreamsWatermarkHoldState(namespace, id, 
timestampCombiner);
+    }
+  }
+
+  /** Common read/write/clear against the backing store for one cell. */
+  private abstract class AbstractState<T> {
+    final StateNamespace namespace;
+    final String id;
+    final Coder<T> coder;
+
+    AbstractState(StateNamespace namespace, String id, Coder<T> coder) {
+      this.namespace = namespace;
+      this.id = id;
+      this.coder = coder;
+    }
+
+    @Nullable
+    T readValue() {
+      byte[] bytes = store.get(compositeKey(namespace, id));
+      if (bytes == null) {
+        return null;
+      }
+      try {
+        return CoderUtils.decodeFromByteArray(coder, bytes);
+      } catch (CoderException e) {
+        throw new RuntimeException("Failed to decode state " + id, e);
+      }
+    }
+
+    void writeValue(T input) {
+      try {
+        store.put(compositeKey(namespace, id), 
CoderUtils.encodeToByteArray(coder, input));
+      } catch (CoderException e) {
+        throw new RuntimeException("Failed to encode state " + id, e);
+      }
+    }
+
+    public void clear() {
+      store.delete(compositeKey(namespace, id));
+    }
+
+    ReadableState<Boolean> isEmptyState() {
+      return new ReadableState<Boolean>() {
+        @Override
+        public Boolean read() {
+          return store.get(compositeKey(namespace, id)) == null;
+        }
+
+        @Override
+        public ReadableState<Boolean> readLater() {
+          return this;
+        }
+      };
+    }
+  }
+
+  private class KafkaStreamsValueState<T> extends AbstractState<T> implements 
ValueState<T> {
+    KafkaStreamsValueState(StateNamespace namespace, String id, Coder<T> 
coder) {
+      super(namespace, id, coder);
+    }
+
+    @Override
+    public KafkaStreamsValueState<T> readLater() {
+      return this;
+    }
+
+    @Override
+    public @Nullable T read() {
+      return readValue();
+    }
+
+    @Override
+    public void write(T input) {
+      writeValue(input);
+    }
+  }
+
+  private class KafkaStreamsBagState<T> extends AbstractState<List<T>> 
implements BagState<T> {
+    KafkaStreamsBagState(StateNamespace namespace, String id, Coder<T> 
elemCoder) {
+      super(namespace, id, ListCoder.of(elemCoder));
+    }
+
+    @Override
+    public KafkaStreamsBagState<T> readLater() {
+      return this;
+    }
+
+    @Override
+    public Iterable<T> read() {
+      List<T> value = readValue();
+      return value == null ? new ArrayList<>() : value;
+    }
+
+    @Override
+    public void add(T input) {
+      List<T> value = readValue();
+      if (value == null) {
+        value = new ArrayList<>();
+      }
+      value.add(input);
+      writeValue(value);
+    }
+
+    @Override
+    public ReadableState<Boolean> isEmpty() {
+      return isEmptyState();
+    }
+  }
+
+  private class KafkaStreamsWatermarkHoldState extends AbstractState<Instant>
+      implements WatermarkHoldState {
+    private final TimestampCombiner timestampCombiner;
+
+    KafkaStreamsWatermarkHoldState(
+        StateNamespace namespace, String id, TimestampCombiner 
timestampCombiner) {
+      super(namespace, id, InstantCoder.of());
+      this.timestampCombiner = timestampCombiner;
+    }
+
+    @Override
+    public KafkaStreamsWatermarkHoldState readLater() {
+      return this;
+    }
+
+    // GroupingState.read() is typed non-null, but an empty hold reads back 
null. Beam's state
+    // interfaces are under-annotated here 
(https://github.com/apache/beam/issues/20497), which is
+    // why the Spark and Flink StateInternals suppress nullness for the whole 
class; this runner
+    // narrows the suppression to just this method.
+    @Override
+    @SuppressWarnings("nullness")
+    public Instant read() {
+      return readValue();
+    }
+
+    @Override
+    public void add(Instant outputTime) {
+      Instant current = readValue();
+      Instant combined =
+          current == null ? outputTime : timestampCombiner.combine(current, 
outputTime);
+      writeValue(combined);
+      // Mirror the hold into the index so the processor can find the minimum 
hold across every key
+      // and window with one lookup instead of reading all of them.
+      if (current != null) {
+        holdsIndexStore.delete(holdIndexKey(current));
+      }
+      holdsIndexStore.put(holdIndexKey(combined), EMPTY_VALUE);
+    }
+
+    @Override
+    public void clear() {
+      Instant current = readValue();
+      if (current != null) {
+        holdsIndexStore.delete(holdIndexKey(current));
+      }
+      super.clear();
+    }
+
+    /** {@code holdTimestamp | cell}, so the index is ordered by hold time. */
+    private byte[] holdIndexKey(Instant hold) {
+      byte[] cellKey = compositeKey(namespace, id);
+      byte[] indexKey = new byte[StoreKeys.TIMESTAMP_BYTES + cellKey.length];
+      int offset = StoreKeys.writeTimestamp(indexKey, 0, hold.getMillis());
+      System.arraycopy(cellKey, 0, indexKey, offset, cellKey.length);
+      return indexKey;
+    }
+
+    @Override
+    public ReadableState<Boolean> isEmpty() {
+      return isEmptyState();
+    }
+
+    @Override
+    public TimestampCombiner getTimestampCombiner() {
+      return timestampCombiner;
+    }
+  }
+
+  @SuppressWarnings("TypeParameterShadowing")
+  private class KafkaStreamsCombiningState<InputT, AccumT, OutputT> extends 
AbstractState<AccumT>
+      implements CombiningState<InputT, AccumT, OutputT> {
+    private final CombineFn<InputT, AccumT, OutputT> combineFn;
+
+    KafkaStreamsCombiningState(
+        StateNamespace namespace,
+        String id,
+        Coder<AccumT> accumCoder,
+        CombineFn<InputT, AccumT, OutputT> combineFn) {
+      super(namespace, id, accumCoder);
+      this.combineFn = combineFn;
+    }
+
+    @Override
+    public KafkaStreamsCombiningState<InputT, AccumT, OutputT> readLater() {
+      return this;
+    }
+
+    // GroupingState.read() is typed non-null but a CombineFn may extract a 
null output; the same
+    // under-annotation as WatermarkHoldState.read() 
(https://github.com/apache/beam/issues/20497).
+    @Override
+    @SuppressWarnings("nullness")
+    public OutputT read() {
+      return combineFn.extractOutput(getAccum());
+    }
+
+    @Override
+    public void add(InputT input) {
+      writeValue(combineFn.addInput(getAccum(), input));
+    }
+
+    @Override
+    public AccumT getAccum() {
+      AccumT accum = readValue();
+      return accum == null ? combineFn.createAccumulator() : accum;
+    }
+
+    @Override
+    public void addAccum(AccumT accum) {
+      writeValue(combineFn.mergeAccumulators(Arrays.asList(getAccum(), 
accum)));
+    }
+
+    @Override
+    public AccumT mergeAccumulators(Iterable<AccumT> accumulators) {
+      return combineFn.mergeAccumulators(accumulators);
+    }
+
+    @Override
+    public ReadableState<Boolean> isEmpty() {
+      return isEmptyState();
+    }
+  }
+}
diff --git 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternals.java
 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternals.java
new file mode 100644
index 00000000000..dddc28eb438
--- /dev/null
+++ 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternals.java
@@ -0,0 +1,265 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.translation;
+
+import java.nio.charset.StandardCharsets;
+import org.apache.beam.runners.core.StateNamespace;
+import org.apache.beam.runners.core.TimerInternals;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.coders.CoderException;
+import org.apache.beam.sdk.state.TimeDomain;
+import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.util.CoderUtils;
+import org.apache.kafka.streams.state.KeyValueStore;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.joda.time.Instant;
+
+/**
+ * A {@link TimerInternals} for one key, backed by two Kafka Streams stores 
shared by a GroupByKey.
+ *
+ * <p>Kafka Streams has no per-key timer service, so timers are persisted like 
any other state, in
+ * two stores that serve the two ways a timer is looked up:
+ *
+ * <ul>
+ *   <li>the <b>identity store</b>, keyed by {@code key | domain | timerFamily 
| timerId |
+ *       namespace}, is how {@link #setTimer} overwrites and {@link 
#deleteTimer} removes exactly
+ *       one timer, as {@link TimerInternals}' contract requires. Its value is 
the index key below,
+ *       so a timer that is overwritten or deleted can have its index entry 
removed without knowing
+ *       what time it had been set for.
+ *   <li>the <b>index store</b>, keyed by {@code domain | fireTimestamp | 
identity}, is how due
+ *       timers are found. Because the timestamp is written in the sortable 
form described on {@link
+ *       StoreKeys}, all event-time timers due at a watermark are one range 
scan — {@link
+ *       #dueEventTimeRangeStart} to {@link #dueEventTimeRangeEnd} — rather 
than a scan of every
+ *       timer of every key. Its value is the {@link TimerData}, so firing 
needs no second lookup.
+ * </ul>
+ *
+ * <p>Firing is driven by {@link WindowedGroupByKeyProcessor}: on a watermark 
advance it range-scans
+ * the index for event-time timers that are due and replays them through {@link
+ * org.apache.beam.runners.core.ReduceFnRunner#onTimers}.
+ *
+ * <p>This instance reports the times it was constructed with; it never fires 
timers itself.
+ */
+class KafkaStreamsTimerInternals implements TimerInternals {
+
+  private final byte[] encodedKey;
+  private final KeyValueStore<byte[], byte[]> identityStore;
+  private final KeyValueStore<byte[], byte[]> indexStore;
+  private final TimerInternals.TimerDataCoderV2 timerCoder;
+  private final Instant inputWatermarkTime;
+  private final Instant outputWatermarkTime;
+  private final Instant processingTime;
+
+  KafkaStreamsTimerInternals(
+      byte[] encodedKey,
+      KeyValueStore<byte[], byte[]> identityStore,
+      KeyValueStore<byte[], byte[]> indexStore,
+      Coder<? extends BoundedWindow> windowCoder,
+      Instant inputWatermarkTime,
+      Instant outputWatermarkTime,
+      Instant processingTime) {
+    this.encodedKey = encodedKey;
+    this.identityStore = identityStore;
+    this.indexStore = indexStore;
+    this.timerCoder = TimerInternals.TimerDataCoderV2.of(windowCoder);
+    this.inputWatermarkTime = inputWatermarkTime;
+    this.outputWatermarkTime = outputWatermarkTime;
+    this.processingTime = processingTime;
+  }
+
+  @Override
+  public void setTimer(
+      StateNamespace namespace,
+      String timerId,
+      String timerFamilyId,
+      Instant target,
+      Instant outputTimestamp,
+      TimeDomain timeDomain) {
+    setTimer(TimerData.of(timerId, timerFamilyId, namespace, target, 
outputTimestamp, timeDomain));
+  }
+
+  @Override
+  public void setTimer(TimerData timerData) {
+    byte[] identityKey = identityKey(encodedKey, timerData);
+    // Setting a timer that already exists replaces it, so drop the old index 
entry first —
+    // otherwise the timer would still be due at the time it was originally 
set for.
+    byte[] previousIndexKey = identityStore.get(identityKey);
+    if (previousIndexKey != null) {
+      indexStore.delete(previousIndexKey);
+    }
+    byte[] indexKey =
+        indexKey(timerData.getDomain(), timerData.getTimestamp().getMillis(), 
identityKey);
+    identityStore.put(identityKey, indexKey);
+    indexStore.put(indexKey, encodeTimer(timerData));
+  }
+
+  @Override
+  public void deleteTimer(
+      StateNamespace namespace, String timerId, String timerFamilyId, 
TimeDomain timeDomain) {
+    deleteByIdentity(identityKey(encodedKey, timerId, timerFamilyId, 
timeDomain, namespace));
+  }
+
+  @Override
+  public void deleteTimer(StateNamespace namespace, String timerId, String 
timerFamilyId) {
+    throw new UnsupportedOperationException(
+        "Deleting a timer without a time domain is not supported; the domain 
is part of a timer's"
+            + " store identity.");
+  }
+
+  @Override
+  public void deleteTimer(TimerData timerKey) {
+    deleteByIdentity(identityKey(encodedKey, timerKey));
+  }
+
+  private void deleteByIdentity(byte[] identityKey) {
+    byte[] indexKey = identityStore.get(identityKey);
+    if (indexKey != null) {
+      indexStore.delete(indexKey);
+    }
+    identityStore.delete(identityKey);
+  }
+
+  @Override
+  public Instant currentProcessingTime() {
+    return processingTime;
+  }
+
+  /**
+   * Returns {@code null}: a synchronized processing time is the slowest 
processing time across the
+   * job's workers, which needs the cross-instance coordination that the 
runner's watermark reports
+   * only carry for event time. {@link TimerInternals} allows null here, and 
nothing on the paths
+   * this runner supports today reads it — it is consulted for processing-time 
triggers, which land
+   * with the processing-time timer support in a follow-up (the same work that 
would supply it).
+   */
+  @Override
+  public @Nullable Instant currentSynchronizedProcessingTime() {
+    return null;
+  }
+
+  @Override
+  public Instant currentInputWatermarkTime() {
+    return inputWatermarkTime;
+  }
+
+  /**
+   * The watermark this GroupByKey has last forwarded downstream, which trails 
{@link
+   * #currentInputWatermarkTime} by the pending watermark holds.
+   */
+  @Override
+  public Instant currentOutputWatermarkTime() {
+    return outputWatermarkTime;
+  }
+
+  private byte[] encodeTimer(TimerData timerData) {
+    try {
+      return CoderUtils.encodeToByteArray(timerCoder, timerData);
+    } catch (CoderException e) {
+      throw new RuntimeException("Failed to encode timer " + timerData, e);
+    }
+  }
+
+  /** Decodes an index store value back into its timer. */
+  static TimerData decodeTimer(Coder<? extends BoundedWindow> windowCoder, 
byte[] bytes) {
+    try {
+      return 
CoderUtils.decodeFromByteArray(TimerInternals.TimerDataCoderV2.of(windowCoder), 
bytes);
+    } catch (CoderException e) {
+      throw new RuntimeException("Failed to decode timer", e);
+    }
+  }
+
+  /**
+   * The identity store key for a timer: {@code key | domain | timerFamily | 
timerId | namespace}.
+   */
+  static byte[] identityKey(byte[] encodedKey, TimerData timerData) {
+    return identityKey(
+        encodedKey,
+        timerData.getTimerId(),
+        timerData.getTimerFamilyId(),
+        timerData.getDomain(),
+        timerData.getNamespace());
+  }
+
+  static byte[] identityKey(
+      byte[] encodedKey,
+      String timerId,
+      String timerFamilyId,
+      TimeDomain domain,
+      StateNamespace namespace) {
+    byte[] domainBytes = {(byte) domain.ordinal()};
+    byte[] familyBytes = timerFamilyId.getBytes(StandardCharsets.UTF_8);
+    byte[] idBytes = timerId.getBytes(StandardCharsets.UTF_8);
+    byte[] namespaceBytes = 
namespace.stringKey().getBytes(StandardCharsets.UTF_8);
+    byte[] key =
+        new byte
+            [StoreKeys.segmentLength(encodedKey)
+                + StoreKeys.segmentLength(domainBytes)
+                + StoreKeys.segmentLength(familyBytes)
+                + StoreKeys.segmentLength(idBytes)
+                + StoreKeys.segmentLength(namespaceBytes)];
+    int offset = StoreKeys.writeSegment(key, 0, encodedKey);
+    offset = StoreKeys.writeSegment(key, offset, domainBytes);
+    offset = StoreKeys.writeSegment(key, offset, familyBytes);
+    offset = StoreKeys.writeSegment(key, offset, idBytes);
+    StoreKeys.writeSegment(key, offset, namespaceBytes);
+    return key;
+  }
+
+  /** The index store key for a timer: {@code domain | fireTimestamp | 
identity}. */
+  static byte[] indexKey(TimeDomain domain, long fireMillis, byte[] 
identityKey) {
+    byte[] key = new byte[1 + StoreKeys.TIMESTAMP_BYTES + identityKey.length];
+    key[0] = (byte) domain.ordinal();
+    int offset = StoreKeys.writeTimestamp(key, 1, fireMillis);
+    System.arraycopy(identityKey, 0, key, offset, identityKey.length);
+    return key;
+  }
+
+  /** Inclusive lower bound of the range scan for due event-time timers. */
+  static byte[] dueEventTimeRangeStart() {
+    byte[] bound = new byte[1 + StoreKeys.TIMESTAMP_BYTES];
+    bound[0] = (byte) TimeDomain.EVENT_TIME.ordinal();
+    StoreKeys.writeTimestamp(bound, 1, Long.MIN_VALUE);
+    return bound;
+  }
+
+  /**
+   * Inclusive upper bound of the range scan for event-time timers due at 
{@code watermarkMillis}.
+   *
+   * <p>Every index key carries a non-empty identity after its timestamp, so 
no key is equal to the
+   * bare {@code domain | watermark + 1} prefix returned here: an inclusive 
scan up to it yields
+   * exactly the timers whose fire time is at or before the watermark. Beam's 
maximum timestamp is
+   * far below {@link Long#MAX_VALUE}, so the increment cannot overflow.
+   */
+  static byte[] dueEventTimeRangeEnd(long watermarkMillis) {
+    byte[] bound = new byte[1 + StoreKeys.TIMESTAMP_BYTES];
+    bound[0] = (byte) TimeDomain.EVENT_TIME.ordinal();
+    StoreKeys.writeTimestamp(bound, 1, watermarkMillis + 1);
+    return bound;
+  }
+
+  /** Reads the identity key back out of an index key. */
+  static byte[] identityKeyOf(byte[] indexKey) {
+    int offset = 1 + StoreKeys.TIMESTAMP_BYTES;
+    byte[] identityKey = new byte[indexKey.length - offset];
+    System.arraycopy(indexKey, offset, identityKey, 0, identityKey.length);
+    return identityKey;
+  }
+
+  /** Reads the encoded Beam key (the first segment) back out of an identity 
key. */
+  static byte[] encodedKeyOf(byte[] identityKey) {
+    return StoreKeys.readSegment(identityKey, 0);
+  }
+}
diff --git 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StoreKeys.java
 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StoreKeys.java
new file mode 100644
index 00000000000..5cb00ca8b6c
--- /dev/null
+++ 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StoreKeys.java
@@ -0,0 +1,105 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.translation;
+
+/**
+ * Byte helpers for the composite keys the runner's state and timer stores are 
addressed by.
+ *
+ * <p>Keys are built into a single pre-sized array rather than through a 
stream, because they are on
+ * the hot path: one is built for every state cell read or written, several 
times per element.
+ *
+ * <p>Variable-length parts are written length-prefixed. A separator byte 
would be shorter, but both
+ * an encoded Beam key (arbitrary user coder output) and a {@link
+ * org.apache.beam.runners.core.StateNamespace#stringKey} can contain any byte 
value, so no
+ * separator is safe from collisions — {@code key="a/b", ns="c"} and {@code 
key="a", ns="b/c"} would
+ * produce the same bytes. Length prefixes also let the encoded key be read 
back out of a timer key,
+ * which the timer scan needs.
+ *
+ * <p>Timestamps are written sign-flipped big-endian so that the unsigned 
lexicographic order Kafka
+ * Streams compares keys by is the same as numeric order. That is what makes a 
range scan over a
+ * timestamp-prefixed store return exactly the entries up to a point in time, 
which is how due
+ * timers and the minimum watermark hold are found without scanning everything.
+ */
+final class StoreKeys {
+
+  /** Bytes taken by a length prefix. */
+  static final int LENGTH_BYTES = 4;
+
+  /** Bytes taken by a sortable timestamp. */
+  static final int TIMESTAMP_BYTES = 8;
+
+  private StoreKeys() {}
+
+  /** Bytes a length-prefixed segment occupies. */
+  static int segmentLength(byte[] segment) {
+    return LENGTH_BYTES + segment.length;
+  }
+
+  /** Writes {@code segment} length-prefixed at {@code offset}, returning the 
offset after it. */
+  static int writeSegment(byte[] target, int offset, byte[] segment) {
+    int next = writeLength(target, offset, segment.length);
+    System.arraycopy(segment, 0, target, next, segment.length);
+    return next + segment.length;
+  }
+
+  private static int writeLength(byte[] target, int offset, int length) {
+    target[offset] = (byte) ((length >>> 24) & 0xff);
+    target[offset + 1] = (byte) ((length >>> 16) & 0xff);
+    target[offset + 2] = (byte) ((length >>> 8) & 0xff);
+    target[offset + 3] = (byte) (length & 0xff);
+    return offset + LENGTH_BYTES;
+  }
+
+  /** Reads the length prefix at {@code offset}. */
+  static int readLength(byte[] source, int offset) {
+    return ((source[offset] & 0xff) << 24)
+        | ((source[offset + 1] & 0xff) << 16)
+        | ((source[offset + 2] & 0xff) << 8)
+        | (source[offset + 3] & 0xff);
+  }
+
+  /** Reads the length-prefixed segment starting at {@code offset}. */
+  static byte[] readSegment(byte[] source, int offset) {
+    int length = readLength(source, offset);
+    byte[] segment = new byte[length];
+    System.arraycopy(source, offset + LENGTH_BYTES, segment, 0, length);
+    return segment;
+  }
+
+  /**
+   * Writes a timestamp so that unsigned byte order matches numeric order: 
flipping the sign bit
+   * maps {@link Long#MIN_VALUE}..{@link Long#MAX_VALUE} onto 
0x00..&nbsp;0xff.. big-endian, so
+   * negative timestamps (valid in Beam) sort before positive ones.
+   */
+  static int writeTimestamp(byte[] target, int offset, long millis) {
+    long sortable = millis ^ Long.MIN_VALUE;
+    for (int i = 0; i < TIMESTAMP_BYTES; i++) {
+      target[offset + i] = (byte) ((sortable >>> (8 * (TIMESTAMP_BYTES - 1 - 
i))) & 0xff);
+    }
+    return offset + TIMESTAMP_BYTES;
+  }
+
+  /** Reads a timestamp written by {@link #writeTimestamp}. */
+  static long readTimestamp(byte[] source, int offset) {
+    long sortable = 0;
+    for (int i = 0; i < TIMESTAMP_BYTES; i++) {
+      sortable = (sortable << 8) | (source[offset + i] & 0xffL);
+    }
+    return sortable ^ Long.MIN_VALUE;
+  }
+}
diff --git 
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WindowedGroupByKeyProcessor.java
 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WindowedGroupByKeyProcessor.java
new file mode 100644
index 00000000000..1b0ffae23fd
--- /dev/null
+++ 
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WindowedGroupByKeyProcessor.java
@@ -0,0 +1,323 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.translation;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.beam.model.pipeline.v1.RunnerApi;
+import org.apache.beam.runners.core.NullSideInputReader;
+import org.apache.beam.runners.core.ReduceFnRunner;
+import org.apache.beam.runners.core.StateInternals;
+import org.apache.beam.runners.core.SystemReduceFn;
+import org.apache.beam.runners.core.TimerInternals.TimerData;
+import org.apache.beam.runners.core.triggers.ExecutableTriggerStateMachine;
+import org.apache.beam.runners.core.triggers.TriggerStateMachines;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.coders.CoderException;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.util.CoderUtils;
+import org.apache.beam.sdk.util.construction.TriggerTranslation;
+import org.apache.beam.sdk.values.KV;
+import org.apache.beam.sdk.values.WindowedValue;
+import org.apache.beam.sdk.values.WindowingStrategy;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.io.BaseEncoding;
+import org.apache.kafka.streams.processor.api.Processor;
+import org.apache.kafka.streams.processor.api.ProcessorContext;
+import org.apache.kafka.streams.processor.api.Record;
+import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.KeyValueStore;
+import org.checkerframework.checker.nullness.qual.NonNull;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.joda.time.Instant;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Executes a windowed {@code GroupByKey} by driving Beam's {@link 
ReduceFnRunner} — the same
+ * windowing + triggering state machine the Flink and Spark portable runners 
use — with Kafka
+ * Streams state and timers behind it.
+ *
+ * <p>Records arrive on the repartition topic keyed by the encoded Beam key, 
so every value of a key
+ * is co-located here. For each data record this builds a {@link 
ReduceFnRunner} for that key over a
+ * {@link KafkaStreamsStateInternals} and {@link KafkaStreamsTimerInternals} 
(both backed by
+ * persistent stores) and feeds the element in; the runner assigns it to 
windows, updates the
+ * trigger state, and sets any timers it needs. When the aggregated input 
watermark advances, every
+ * event-time timer whose fire time has passed is replayed through {@code 
onTimers}, which is what
+ * makes windows emit their panes. The runner is stateless between records — 
all durable state lives
+ * in the stores — so a fresh one per record is correct, mirroring {@code
+ * GroupAlsoByWindowViaWindowSetNewDoFn}.
+ *
+ * <p>This first version supports the default trigger and non-merging windows 
well; richer triggers,
+ * processing-time timers and session (merging) windows build on the same 
machinery in a follow-up.
+ */
+class WindowedGroupByKeyProcessor<K, V, W extends BoundedWindow>
+    implements Processor<byte[], KStreamsPayload<?>, byte[], 
KStreamsPayload<?>> {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(WindowedGroupByKeyProcessor.class);
+
+  private final String stateStoreName;
+  private final String holdsIndexStoreName;
+  private final String timerStoreName;
+  private final String timerIndexStoreName;
+  private final String transformId;
+  private final Coder<K> keyCoder;
+  private final WindowingStrategy<?, W> windowingStrategy;
+  private final Coder<? extends BoundedWindow> windowCoder;
+  private final RunnerApi.Trigger triggerProto;
+  private final SystemReduceFn<K, V, Iterable<V>, Iterable<V>, W> reduceFn;
+  private final PipelineOptions options;
+
+  private final WatermarkAggregator watermarkAggregator;
+  private Instant lastForwardedWatermark = BoundedWindow.TIMESTAMP_MIN_VALUE;
+  private Instant inputWatermark = BoundedWindow.TIMESTAMP_MIN_VALUE;
+
+  private @Nullable ProcessorContext<byte[], KStreamsPayload<?>> context;
+  private @Nullable KeyValueStore<byte[], byte[]> stateStore;
+  private @Nullable KeyValueStore<byte[], byte[]> holdsIndexStore;
+  private @Nullable KeyValueStore<byte[], byte[]> timerStore;
+  private @Nullable KeyValueStore<byte[], byte[]> timerIndexStore;
+
+  WindowedGroupByKeyProcessor(
+      String stateStoreName,
+      String holdsIndexStoreName,
+      String timerStoreName,
+      String timerIndexStoreName,
+      String transformId,
+      Set<String> upstreamTransformIds,
+      Coder<K> keyCoder,
+      Coder<V> valueCoder,
+      WindowingStrategy<?, W> windowingStrategy,
+      PipelineOptions options) {
+    this.stateStoreName = stateStoreName;
+    this.holdsIndexStoreName = holdsIndexStoreName;
+    this.timerStoreName = timerStoreName;
+    this.timerIndexStoreName = timerIndexStoreName;
+    this.transformId = transformId;
+    this.keyCoder = keyCoder;
+    this.windowingStrategy = windowingStrategy;
+    this.windowCoder = windowingStrategy.getWindowFn().windowCoder();
+    this.triggerProto = 
TriggerTranslation.toProto(windowingStrategy.getTrigger());
+    this.reduceFn = SystemReduceFn.buffering(valueCoder);
+    this.options = options;
+    this.watermarkAggregator = new WatermarkAggregator(upstreamTransformIds);
+  }
+
+  @Override
+  public void init(ProcessorContext<byte[], KStreamsPayload<?>> context) {
+    this.context = context;
+    this.stateStore = context.getStateStore(stateStoreName);
+    this.holdsIndexStore = context.getStateStore(holdsIndexStoreName);
+    this.timerStore = context.getStateStore(timerStoreName);
+    this.timerIndexStore = context.getStateStore(timerIndexStoreName);
+  }
+
+  @Override
+  public void process(Record<byte[], KStreamsPayload<?>> record) {
+    KStreamsPayload<?> payload = record.value();
+    if (payload == null) {
+      LOG.warn(
+          "GroupByKey {} dropping record with null payload (external write or 
tombstone)",
+          transformId);
+      return;
+    }
+    if (payload.isData()) {
+      processData(record, payload);
+      return;
+    }
+    watermarkAggregator.observe(payload.asWatermark());
+    Instant advanced = watermarkAggregator.advance();
+    if (advanced.isAfter(inputWatermark)) {
+      inputWatermark = advanced;
+      fireDueEventTimeTimers(record, advanced);
+    }
+    // Firing may have emitted panes and released their holds, so the output 
watermark is computed
+    // after it.
+    Instant output = outputWatermark();
+    if (output.isAfter(lastForwardedWatermark)) {
+      lastForwardedWatermark = output;
+      forwardWatermark(record, output.getMillis());
+    }
+  }
+
+  /**
+   * The watermark to publish downstream: the input watermark, held back by 
the earliest watermark
+   * hold any pending pane has taken.
+   *
+   * <p>{@link ReduceFnRunner} takes a hold for buffered elements that have 
not been emitted yet, at
+   * the timestamp their pane will carry. Forwarding the raw input watermark 
would tell downstream
+   * that nothing earlier is coming while those panes are still buffered, and 
the elements would
+   * then arrive late against the watermark we had already published.
+   */
+  private Instant outputWatermark() {
+    Instant minHold =
+        
KafkaStreamsStateInternals.minWatermarkHold(checkInitialized(holdsIndexStore));
+    return minHold == null || inputWatermark.isBefore(minHold) ? 
inputWatermark : minHold;
+  }
+
+  private void processData(Record<byte[], KStreamsPayload<?>> record, 
KStreamsPayload<?> payload) {
+    byte[] encodedKey = record.key();
+    if (encodedKey == null) {
+      throw new IllegalStateException("GroupByKey data record is missing its 
key");
+    }
+    @SuppressWarnings("unchecked")
+    WindowedValue<KV<K, V>> element = (WindowedValue<KV<K, V>>) 
payload.getData();
+    K key = decodeKey(encodedKey);
+    WindowedValue<V> valueElement = 
element.withValue(element.getValue().getValue());
+    runReduceFn(
+        record, encodedKey, key, Collections.singletonList(valueElement), 
Collections.emptyList());
+  }
+
+  /**
+   * Fires every event-time timer whose fire time is at or before the new 
input watermark.
+   *
+   * <p>The timers to fire are found by range-scanning the fire-time-ordered 
index over exactly the
+   * window {@code (-inf, watermark]}, so the cost is proportional to the 
number of timers that are
+   * actually due rather than to the number of keys that hold a timer.
+   */
+  private void fireDueEventTimeTimers(
+      Record<byte[], KStreamsPayload<?>> record, Instant watermark) {
+    KeyValueStore<byte[], byte[]> identityStore = checkInitialized(timerStore);
+    KeyValueStore<byte[], byte[]> indexStore = 
checkInitialized(timerIndexStore);
+    // Group the due timers by the Beam key they belong to; a key's timers 
fire together in one
+    // ReduceFnRunner turn.
+    Map<String, DueTimers> dueByKey = new LinkedHashMap<>();
+    List<byte[]> firedIndexKeys = new ArrayList<>();
+    try (KeyValueIterator<byte[], byte[]> it =
+        indexStore.range(
+            KafkaStreamsTimerInternals.dueEventTimeRangeStart(),
+            
KafkaStreamsTimerInternals.dueEventTimeRangeEnd(watermark.getMillis()))) {
+      while (it.hasNext()) {
+        org.apache.kafka.streams.KeyValue<byte[], byte[]> entry = it.next();
+        TimerData timer = KafkaStreamsTimerInternals.decodeTimer(windowCoder, 
entry.value);
+        byte[] encodedKey =
+            KafkaStreamsTimerInternals.encodedKeyOf(
+                KafkaStreamsTimerInternals.identityKeyOf(entry.key));
+        dueByKey
+            .computeIfAbsent(
+                BaseEncoding.base16().encode(encodedKey), k -> new 
DueTimers(encodedKey))
+            .timers
+            .add(timer);
+        firedIndexKeys.add(entry.key);
+      }
+    }
+    // Clear the fired timers from both stores before replaying them, since 
onTimers may
+    // legitimately
+    // set new ones — including at the same identity.
+    for (byte[] indexKey : firedIndexKeys) {
+      indexStore.delete(indexKey);
+      identityStore.delete(KafkaStreamsTimerInternals.identityKeyOf(indexKey));
+    }
+    for (DueTimers due : dueByKey.values()) {
+      runReduceFn(
+          record, due.encodedKey, decodeKey(due.encodedKey), 
Collections.emptyList(), due.timers);
+    }
+  }
+
+  private void runReduceFn(
+      Record<byte[], KStreamsPayload<?>> record,
+      byte[] encodedKey,
+      @NonNull K key,
+      List<WindowedValue<V>> elements,
+      List<TimerData> timers) {
+    StateInternals stateInternals =
+        new KafkaStreamsStateInternals<>(
+            key, encodedKey, checkInitialized(stateStore), 
checkInitialized(holdsIndexStore));
+    KafkaStreamsTimerInternals timerInternals =
+        new KafkaStreamsTimerInternals(
+            encodedKey,
+            checkInitialized(timerStore),
+            checkInitialized(timerIndexStore),
+            windowCoder,
+            inputWatermark,
+            lastForwardedWatermark,
+            Instant.now());
+    ReduceFnRunner<K, V, Iterable<V>, W> runner =
+        new ReduceFnRunner<>(
+            key,
+            windowingStrategy,
+            ExecutableTriggerStateMachine.create(
+                TriggerStateMachines.stateMachineForTrigger(triggerProto)),
+            stateInternals,
+            timerInternals,
+            output -> forwardData(record, encodedKey, output),
+            NullSideInputReader.empty(),
+            reduceFn,
+            options);
+    try {
+      runner.processElements(elements);
+      runner.onTimers(timers);
+      runner.persist();
+    } catch (Exception e) {
+      throw new RuntimeException("GroupByKey " + transformId + " failed to run 
windowing", e);
+    }
+  }
+
+  private void forwardData(
+      Record<byte[], KStreamsPayload<?>> trigger,
+      byte[] encodedKey,
+      WindowedValue<KV<K, Iterable<V>>> output) {
+    ProcessorContext<byte[], KStreamsPayload<?>> ctx = 
checkInitialized(context);
+    ctx.forward(
+        new Record<byte[], KStreamsPayload<?>>(
+            encodedKey, KStreamsPayload.data(output), trigger.timestamp()));
+  }
+
+  private void forwardWatermark(Record<byte[], KStreamsPayload<?>> trigger, 
long watermarkMillis) {
+    ProcessorContext<byte[], KStreamsPayload<?>> ctx = 
checkInitialized(context);
+    // Stamped with this transform's own id; GroupByKey is a single instance 
for now (0 of 1).
+    ctx.forward(
+        new Record<byte[], KStreamsPayload<?>>(
+            trigger.key(),
+            KStreamsPayload.watermark(watermarkMillis, transformId, 0, 1),
+            trigger.timestamp()));
+  }
+
+  private @NonNull K decodeKey(byte[] bytes) {
+    try {
+      K key = CoderUtils.decodeFromByteArray(keyCoder, bytes);
+      if (key == null) {
+        throw new IllegalStateException("GroupByKey key decoded to null");
+      }
+      return key;
+    } catch (CoderException e) {
+      throw new RuntimeException("Failed to decode GroupByKey key", e);
+    }
+  }
+
+  private static <T> T checkInitialized(@Nullable T value) {
+    if (value == null) {
+      throw new IllegalStateException("WindowedGroupByKeyProcessor used before 
init()");
+    }
+    return value;
+  }
+
+  /** The event-time timers due for one Beam key, plus that key's encoded 
bytes. */
+  private static final class DueTimers {
+    final byte[] encodedKey;
+    final List<TimerData> timers = new ArrayList<>();
+
+    DueTimers(byte[] encodedKey) {
+      this.encodedKey = encodedKey;
+    }
+  }
+}
diff --git 
a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTestRunner.java
 
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTestRunner.java
index 8ab29182c28..bdbde1db36d 100644
--- 
a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTestRunner.java
+++ 
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTestRunner.java
@@ -18,11 +18,7 @@
 package org.apache.beam.runners.kafka.streams;
 
 import java.time.Duration;
-import java.util.ArrayList;
-import java.util.HashSet;
-import java.util.List;
 import java.util.Properties;
-import java.util.Set;
 import java.util.UUID;
 import org.apache.beam.model.pipeline.v1.RunnerApi;
 import org.apache.beam.runners.core.metrics.MetricsContainerStepMap;
@@ -39,16 +35,11 @@ import org.apache.beam.sdk.util.construction.Environments;
 import org.apache.beam.sdk.util.construction.PipelineOptionsTranslation;
 import org.apache.beam.sdk.util.construction.PipelineTranslation;
 import org.apache.beam.sdk.util.construction.SplittableParDo;
-import org.apache.kafka.common.serialization.ByteArrayDeserializer;
-import org.apache.kafka.common.serialization.ByteArraySerializer;
 import org.apache.kafka.common.serialization.Serdes;
 import org.apache.kafka.streams.StreamsConfig;
-import org.apache.kafka.streams.TestInputTopic;
-import org.apache.kafka.streams.TestOutputTopic;
 import org.apache.kafka.streams.Topology;
 import org.apache.kafka.streams.TopologyDescription;
 import org.apache.kafka.streams.TopologyTestDriver;
-import org.apache.kafka.streams.test.TestRecord;
 
 /**
  * Test harness that runs a Beam {@link Pipeline} through the Kafka Streams 
runner's translation and
@@ -58,16 +49,13 @@ import org.apache.kafka.streams.test.TestRecord;
  * effects (e.g. a {@code SharedTestCollector} written by a recording DoFn) 
have completed when it
  * returns.
  *
- * <p>{@link TopologyTestDriver} does not loop a low-level sink topic back 
into its source, so an
- * internal repartition topic (one that is both a sink and a source in the 
topology — e.g. the one
- * GroupByKey introduces) would otherwise dead-end. {@link #run(Pipeline)} 
discovers those topics
- * from the {@link TopologyDescription} and round-trips them until no more 
records flow, standing in
- * for the broker.
+ * <p>{@link TopologyTestDriver} loops each internal repartition topic (one 
that is both a sink and
+ * a source in the topology — e.g. the one GroupByKey introduces) from its 
sink back to its source
+ * within a single driver step, so advancing the wall clock is enough to drive 
the whole pipeline to
+ * completion; no manual broker simulation is needed.
  */
 public final class KafkaStreamsTestRunner {
 
-  private static final int MAX_ROUND_TRIPS = 100;
-
   private KafkaStreamsTestRunner() {}
 
   /** Pipeline options for a Kafka Streams runner test: the EMBEDDED harness 
and a unique app id. */
@@ -118,10 +106,12 @@ public final class KafkaStreamsTestRunner {
     KafkaStreamsTranslationContext context = translate(pipeline);
     Topology topology = context.getTopology();
     try (TopologyTestDriver driver = new TopologyTestDriver(topology, 
streamsConfig(pipeline))) {
-      // Fire the Impulse wall-clock punctuator and let the initial records 
flow.
+      // Fire the Impulse wall-clock punctuator; TopologyTestDriver then flows 
the records through
+      // the whole topology, including looping each internal repartition topic 
(a sink that is also
+      // a source, e.g. the one GroupByKey introduces) back to its source, 
standing in for the
+      // broker. A second advance covers punctuators that need a later tick.
       driver.advanceWallClockTime(Duration.ofSeconds(1));
       driver.advanceWallClockTime(Duration.ofSeconds(1));
-      roundTripInternalTopics(driver, internalTopics(topology));
     }
     return MetricsContainerStepMap.asAttemptedOnlyMetricResults(
         context.getMetricsContainerStepMap());
@@ -143,81 +133,6 @@ public final class KafkaStreamsTestRunner {
     throw new IllegalStateException("no leaf processor found in topology");
   }
 
-  /** Repartition/internal topics are the ones that appear as both a sink and 
a source. */
-  private static Set<String> internalTopics(Topology topology) {
-    Set<String> sinkTopics = new HashSet<>();
-    Set<String> sourceTopics = new HashSet<>();
-    for (TopologyDescription.Subtopology subtopology : 
topology.describe().subtopologies()) {
-      for (TopologyDescription.Node node : subtopology.nodes()) {
-        if (node instanceof TopologyDescription.Sink) {
-          String topic = ((TopologyDescription.Sink) node).topic();
-          if (topic != null) {
-            sinkTopics.add(topic);
-          }
-        } else if (node instanceof TopologyDescription.Source) {
-          sourceTopics.addAll(((TopologyDescription.Source) node).topicSet());
-        }
-      }
-    }
-    sinkTopics.retainAll(sourceTopics);
-    return sinkTopics;
-  }
-
-  /**
-   * Simulates the broker for internal repartition topics.
-   *
-   * <p>The runner shuffles data (and the watermark) through internal topics 
that a processor both
-   * writes to (a sink) and reads back from (a source) — e.g. the topic 
GroupByKey introduces to
-   * partition by key. On a real broker those records make the round trip 
automatically, but {@link
-   * TopologyTestDriver} does not connect a sink back to a source, so the 
downstream half of the
-   * topology would never see them. This drains what each internal topic's 
sink wrote and pipes it
-   * into that topic's source, repeating until nothing new flows (a fixpoint), 
which stands in for
-   * the broker and lets the pipeline run to completion.
-   */
-  private static void roundTripInternalTopics(TopologyTestDriver driver, 
Set<String> topics) {
-    // Create the sink-output and source-input handles once and reuse them 
across rounds; a single
-    // TestOutputTopic keeps returning newly produced records on each read.
-    List<TopicRoundTrip> roundTrips = new ArrayList<>();
-    for (String topic : topics) {
-      roundTrips.add(
-          new TopicRoundTrip(
-              driver.createOutputTopic(
-                  topic, new ByteArrayDeserializer(), new 
ByteArrayDeserializer()),
-              driver.createInputTopic(
-                  topic, new ByteArraySerializer(), new 
ByteArraySerializer())));
-    }
-
-    for (int round = 0; round < MAX_ROUND_TRIPS; round++) {
-      boolean progressed = false;
-      for (TopicRoundTrip roundTrip : roundTrips) {
-        List<TestRecord<byte[], byte[]>> records = 
roundTrip.output.readRecordsToList();
-        if (records.isEmpty()) {
-          continue;
-        }
-        progressed = true;
-        for (TestRecord<byte[], byte[]> record : records) {
-          roundTrip.input.pipeInput(record);
-        }
-      }
-      if (!progressed) {
-        return;
-      }
-    }
-    throw new IllegalStateException(
-        "Internal topics did not reach quiescence after " + MAX_ROUND_TRIPS + 
" round trips");
-  }
-
-  /** The reusable sink-output and source-input handles for one internal 
topic. */
-  private static final class TopicRoundTrip {
-    final TestOutputTopic<byte[], byte[]> output;
-    final TestInputTopic<byte[], byte[]> input;
-
-    TopicRoundTrip(TestOutputTopic<byte[], byte[]> output, 
TestInputTopic<byte[], byte[]> input) {
-      this.output = output;
-      this.input = input;
-    }
-  }
-
   /** Kafka Streams config for a {@link TopologyTestDriver} built from the 
pipeline's app id. */
   public static Properties streamsConfig(Pipeline pipeline) {
     String applicationId =
diff --git 
a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FixedWindowGroupByKeyTest.java
 
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FixedWindowGroupByKeyTest.java
new file mode 100644
index 00000000000..73f4dac601e
--- /dev/null
+++ 
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FixedWindowGroupByKeyTest.java
@@ -0,0 +1,106 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.translation;
+
+import static org.hamcrest.CoreMatchers.hasItems;
+import static org.hamcrest.CoreMatchers.is;
+import static org.hamcrest.MatcherAssert.assertThat;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import org.apache.beam.runners.kafka.streams.KafkaStreamsTestRunner;
+import org.apache.beam.sdk.Pipeline;
+import org.apache.beam.sdk.coders.KvCoder;
+import org.apache.beam.sdk.coders.StringUtf8Coder;
+import org.apache.beam.sdk.coders.VarIntCoder;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.GroupByKey;
+import org.apache.beam.sdk.transforms.Impulse;
+import org.apache.beam.sdk.transforms.ParDo;
+import org.apache.beam.sdk.transforms.windowing.FixedWindows;
+import org.apache.beam.sdk.transforms.windowing.Window;
+import org.apache.beam.sdk.values.KV;
+import org.joda.time.Duration;
+import org.joda.time.Instant;
+import org.junit.Test;
+
+/**
+ * End-to-end test that GroupByKey groups per fixed window, not just per key: 
{@code Impulse -> emit
+ * timestamped KVs -> Window.into(FixedWindows) -> GroupByKey -> record 
groups}.
+ *
+ * <p>The same key "a" has values in two different windows, so a correct 
windowed GroupByKey emits
+ * two groups for it (one per window) rather than one combined group. This 
exercises the {@link
+ * WindowedGroupByKeyProcessor} path (ReduceFnRunner over the Kafka Streams 
state and timer stores)
+ * that the earlier global-window GroupByKey did not.
+ */
+public class FixedWindowGroupByKeyTest {
+
+  private static final Duration WINDOW_SIZE = Duration.millis(10);
+
+  /** Emits KVs whose timestamps fall into two adjacent fixed windows. */
+  private static class EmitTimestampedKvsFn extends DoFn<byte[], KV<String, 
Integer>> {
+    @ProcessElement
+    public void processElement(OutputReceiver<KV<String, Integer>> out) {
+      // Window [0, 10): a=1, a=2, b=5.
+      out.outputWithTimestamp(KV.of("a", 1), new Instant(1));
+      out.outputWithTimestamp(KV.of("a", 2), new Instant(2));
+      out.outputWithTimestamp(KV.of("b", 5), new Instant(3));
+      // Window [10, 20): a=3.
+      out.outputWithTimestamp(KV.of("a", 3), new Instant(15));
+    }
+  }
+
+  /** Records each grouped result as {@code "key=[sorted values]"}. */
+  private static class RecordGroupFn extends DoFn<KV<String, 
Iterable<Integer>>, Void> {
+    private final SharedTestCollector<String> collector;
+
+    RecordGroupFn(SharedTestCollector<String> collector) {
+      this.collector = collector;
+    }
+
+    @ProcessElement
+    public void processElement(@Element KV<String, Iterable<Integer>> group) {
+      List<Integer> values = new ArrayList<>();
+      group.getValue().forEach(values::add);
+      Collections.sort(values);
+      collector.record(group.getKey() + "=" + values);
+    }
+  }
+
+  @Test
+  public void groupsValuesPerFixedWindow() {
+    try (SharedTestCollector<String> collector = SharedTestCollector.create()) 
{
+      Pipeline pipeline = 
Pipeline.create(KafkaStreamsTestRunner.testOptions());
+      pipeline
+          .apply("impulse", Impulse.create())
+          .apply("emit", ParDo.of(new EmitTimestampedKvsFn()))
+          .setCoder(KvCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))
+          .apply("window", Window.into(FixedWindows.of(WINDOW_SIZE)))
+          .apply("gbk", GroupByKey.create())
+          .apply("record", ParDo.of(new RecordGroupFn(collector)));
+
+      KafkaStreamsTestRunner.run(pipeline);
+
+      List<String> groups = collector.recorded();
+      // a splits across two windows -> two groups; b has one; three groups 
total.
+      assertThat(groups.size(), is(3));
+      assertThat(groups, hasItems("a=[1, 2]", "a=[3]", "b=[5]"));
+    }
+  }
+}
diff --git 
a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternalsTest.java
 
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternalsTest.java
new file mode 100644
index 00000000000..df5b0eda1b2
--- /dev/null
+++ 
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternalsTest.java
@@ -0,0 +1,208 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.translation;
+
+import static org.hamcrest.CoreMatchers.is;
+import static org.hamcrest.MatcherAssert.assertThat;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.beam.runners.core.StateNamespace;
+import org.apache.beam.runners.core.StateNamespaces;
+import org.apache.beam.runners.core.TimerInternals.TimerData;
+import org.apache.beam.sdk.coders.StringUtf8Coder;
+import org.apache.beam.sdk.state.TimeDomain;
+import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.transforms.windowing.GlobalWindow;
+import org.apache.beam.sdk.util.CoderUtils;
+import org.apache.kafka.common.serialization.Serdes;
+import org.apache.kafka.streams.processor.api.MockProcessorContext;
+import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.KeyValueStore;
+import org.apache.kafka.streams.state.Stores;
+import org.joda.time.Instant;
+import org.junit.Before;
+import org.junit.Test;
+
+/**
+ * Tests the timer store and its fire-time index: that a timer can be set, 
replaced and deleted by
+ * identity, and that the timers due at a watermark are found by a range scan 
over the index rather
+ * than by inspecting every timer.
+ */
+public class KafkaStreamsTimerInternalsTest {
+
+  private static final StateNamespace NAMESPACE =
+      StateNamespaces.window(GlobalWindow.Coder.INSTANCE, 
GlobalWindow.INSTANCE);
+
+  private KeyValueStore<byte[], byte[]> identityStore;
+  private KeyValueStore<byte[], byte[]> indexStore;
+
+  @Before
+  public void setUp() {
+    MockProcessorContext<Void, Void> context = new MockProcessorContext<>();
+    identityStore = newStore("timers", context);
+    indexStore = newStore("timers-index", context);
+  }
+
+  private static KeyValueStore<byte[], byte[]> newStore(
+      String name, MockProcessorContext<Void, Void> context) {
+    KeyValueStore<byte[], byte[]> store =
+        Stores.keyValueStoreBuilder(
+                Stores.inMemoryKeyValueStore(name), Serdes.ByteArray(), 
Serdes.ByteArray())
+            .withLoggingDisabled()
+            .build();
+    store.init(context.getStateStoreContext(), store);
+    return store;
+  }
+
+  private KafkaStreamsTimerInternals timersFor(String key) {
+    return new KafkaStreamsTimerInternals(
+        encode(key),
+        identityStore,
+        indexStore,
+        GlobalWindow.Coder.INSTANCE,
+        BoundedWindow.TIMESTAMP_MIN_VALUE,
+        BoundedWindow.TIMESTAMP_MIN_VALUE,
+        new Instant(0));
+  }
+
+  private static byte[] encode(String key) {
+    try {
+      return CoderUtils.encodeToByteArray(StringUtf8Coder.of(), key);
+    } catch (Exception e) {
+      throw new RuntimeException(e);
+    }
+  }
+
+  private static TimerData eventTimer(String id, long millis) {
+    return TimerData.of(
+        id, "", NAMESPACE, new Instant(millis), new Instant(millis), 
TimeDomain.EVENT_TIME);
+  }
+
+  /** The timers the processor would fire at {@code watermarkMillis}, in 
fire-time order. */
+  private List<TimerData> dueAt(long watermarkMillis) {
+    List<TimerData> due = new ArrayList<>();
+    try (KeyValueIterator<byte[], byte[]> it =
+        indexStore.range(
+            KafkaStreamsTimerInternals.dueEventTimeRangeStart(),
+            KafkaStreamsTimerInternals.dueEventTimeRangeEnd(watermarkMillis))) 
{
+      while (it.hasNext()) {
+        due.add(
+            
KafkaStreamsTimerInternals.decodeTimer(GlobalWindow.Coder.INSTANCE, 
it.next().value));
+      }
+    }
+    return due;
+  }
+
+  private static int storeSize(KeyValueStore<byte[], byte[]> store) {
+    int size = 0;
+    try (KeyValueIterator<byte[], byte[]> it = store.all()) {
+      while (it.hasNext()) {
+        it.next();
+        size++;
+      }
+    }
+    return size;
+  }
+
+  @Test
+  public void dueScanReturnsOnlyTimersAtOrBeforeTheWatermark() {
+    KafkaStreamsTimerInternals timers = timersFor("key");
+    timers.setTimer(eventTimer("early", 100L));
+    timers.setTimer(eventTimer("onWatermark", 200L));
+    timers.setTimer(eventTimer("late", 300L));
+
+    List<TimerData> due = dueAt(200L);
+
+    // Ordered by fire time, and the timer set exactly at the watermark is 
included.
+    assertThat(due.size(), is(2));
+    assertThat(due.get(0).getTimerId(), is("early"));
+    assertThat(due.get(1).getTimerId(), is("onWatermark"));
+  }
+
+  @Test
+  public void negativeTimestampsSortBeforePositiveOnes() {
+    KafkaStreamsTimerInternals timers = timersFor("key");
+    timers.setTimer(eventTimer("negative", -5000L));
+    timers.setTimer(eventTimer("zero", 0L));
+    timers.setTimer(eventTimer("positive", 5000L));
+
+    List<TimerData> due = dueAt(0L);
+
+    assertThat(due.size(), is(2));
+    assertThat(due.get(0).getTimerId(), is("negative"));
+    assertThat(due.get(1).getTimerId(), is("zero"));
+  }
+
+  @Test
+  public void resettingATimerReplacesItsIndexEntry() {
+    KafkaStreamsTimerInternals timers = timersFor("key");
+    timers.setTimer(eventTimer("timer", 100L));
+    // Re-setting the same timer identity for a later time must not leave the 
old entry behind,
+    // or the timer would still fire at the time it was first set for.
+    timers.setTimer(eventTimer("timer", 900L));
+
+    assertThat(dueAt(100L).isEmpty(), is(true));
+    assertThat(dueAt(900L).size(), is(1));
+    assertThat(storeSize(indexStore), is(1));
+    assertThat(storeSize(identityStore), is(1));
+  }
+
+  @Test
+  public void deletingATimerRemovesItFromBothStores() {
+    KafkaStreamsTimerInternals timers = timersFor("key");
+    timers.setTimer(eventTimer("timer", 100L));
+    timers.deleteTimer(NAMESPACE, "timer", "", TimeDomain.EVENT_TIME);
+
+    assertThat(dueAt(1000L).isEmpty(), is(true));
+    assertThat(storeSize(indexStore), is(0));
+    assertThat(storeSize(identityStore), is(0));
+  }
+
+  @Test
+  public void timersOfDifferentKeysAreIndependentButShareTheIndex() {
+    timersFor("a").setTimer(eventTimer("timer", 100L));
+    timersFor("b").setTimer(eventTimer("timer", 150L));
+
+    // Same timer id under two Beam keys are two distinct timers, and one scan 
finds both.
+    assertThat(storeSize(identityStore), is(2));
+    assertThat(dueAt(200L).size(), is(2));
+
+    timersFor("a").deleteTimer(NAMESPACE, "timer", "", TimeDomain.EVENT_TIME);
+    assertThat(dueAt(200L).size(), is(1));
+  }
+
+  @Test
+  public void processingTimeTimersAreNotReturnedByTheEventTimeScan() {
+    KafkaStreamsTimerInternals timers = timersFor("key");
+    timers.setTimer(
+        TimerData.of(
+            "processing",
+            "",
+            NAMESPACE,
+            new Instant(100L),
+            new Instant(100L),
+            TimeDomain.PROCESSING_TIME));
+    timers.setTimer(eventTimer("event", 100L));
+
+    List<TimerData> due = dueAt(1000L);
+
+    assertThat(due.size(), is(1));
+    assertThat(due.get(0).getTimerId(), is("event"));
+  }
+}

Reply via email to