gemini-code-assist[bot] commented on code in PR #39410:
URL: https://github.com/apache/beam/pull/39410#discussion_r3623575798


##########
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java:
##########
@@ -249,12 +268,20 @@ private void closeBundleAndFlush(Record<byte[], 
KStreamsPayload<?>> record) {
     }
     ProcessorContext<byte[], KStreamsPayload<?>> ctx = 
checkInitialized(context);
     // The harness has finished the bundle (close() returned) so no further 
enqueues happen.
-    // Drain via poll() so each element is removed as it is forwarded.
-    WindowedValue<?> output;
+    // Drain via poll() so each element is removed as it is forwarded. Each 
output is routed to its
+    // own output's relay child for a multi-output stage; a single-output 
stage forwards directly to
+    // its one downstream (empty routing map).
+    PendingOutput output;
     while ((output = pendingOutputs.poll()) != null) {
-      ctx.forward(
+      Record<byte[], KStreamsPayload<?>> outputRecord =
           new Record<byte[], KStreamsPayload<?>>(
-              record.key(), KStreamsPayload.data(output), record.timestamp()));
+              record.key(), KStreamsPayload.data(output.value), 
record.timestamp());

Review Comment:
   ![medium](https://www.gstatic.com/codereviewagent/medium-priority.svg)
   
   Using `record.withValue(...)` is more concise and idiomatic than manually 
reconstructing the `Record` object. More importantly, it automatically 
preserves any metadata headers present on the original input record, which is 
crucial for features like distributed tracing (e.g., OpenTelemetry/Jaeger) or 
custom interceptors.
   
   ```java
         Record<byte[], KStreamsPayload<?>> outputRecord =
             record.withValue(KStreamsPayload.data(output.value));
   ```



##########
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StageOutputProcessor.java:
##########
@@ -0,0 +1,88 @@
+/*
+ * 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 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.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * One output port of a multi-output {@link ExecutableStageProcessor}: a relay 
node that stands in
+ * as the producer of one of the stage's output PCollections.
+ *
+ * <p>A Kafka Streams processor forwards a record either to all of its 
children or to one named
+ * child, but downstream transforms are wired to a <em>producer node</em> by 
PCollection id, and
+ * that node's identity has to be known when the stage is translated — before 
the downstream
+ * transforms are. So each output PCollection of a multi-output stage gets its 
own relay: the stage
+ * routes each harness output to the matching relay by name, and downstream 
transforms wire to the
+ * relay. (A single-output stage needs none of this and forwards directly.)
+ *
+ * <p>The relay forwards data records unchanged, and re-stamps the watermark 
with its own transform
+ * id so a downstream watermark aggregator, which expects to hear from this 
relay, sees a consistent
+ * source. The stage is a single instance and its watermark is already the 
monotonic aggregate over
+ * its input, so the relay only relabels it — there is exactly one upstream 
partition to relay, no
+ * aggregation to do.
+ */
+class StageOutputProcessor
+    implements Processor<byte[], KStreamsPayload<?>, byte[], 
KStreamsPayload<?>> {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(StageOutputProcessor.class);
+
+  private final String transformId;
+  private @Nullable ProcessorContext<byte[], KStreamsPayload<?>> context;
+
+  StageOutputProcessor(String transformId) {
+    this.transformId = transformId;
+  }
+
+  @Override
+  public void init(ProcessorContext<byte[], KStreamsPayload<?>> context) {
+    this.context = context;
+  }
+
+  @Override
+  public void process(Record<byte[], KStreamsPayload<?>> record) {
+    KStreamsPayload<?> payload = record.value();
+    ProcessorContext<byte[], KStreamsPayload<?>> ctx = context;
+    if (ctx == null) {
+      throw new IllegalStateException("StageOutputProcessor used before 
init()");
+    }
+    if (payload == null) {
+      LOG.warn(
+          "Stage output {} dropping record with null payload (external write 
or tombstone)",
+          transformId);
+      return;
+    }
+    if (!payload.isWatermark()) {
+      // Data for this output: forward unchanged.
+      ctx.forward(record);
+      return;
+    }
+    // Re-stamp the stage's watermark with this output port's own id. Single 
upstream instance, so
+    // its value is already monotonic; just relabel it as source 0 of 1.
+    WatermarkPayload report = payload.asWatermark();
+    ctx.forward(
+        new Record<byte[], KStreamsPayload<?>>(
+            record.key(),
+            KStreamsPayload.watermark(report.getWatermarkMillis(), 
transformId, 0, 1),
+            record.timestamp()));

Review Comment:
   ![medium](https://www.gstatic.com/codereviewagent/medium-priority.svg)
   
   Using `record.withValue(...)` is more concise and idiomatic than manually 
reconstructing the `Record` object. It also ensures that any metadata headers 
present on the original record are preserved when forwarding the watermark 
payload.
   
   ```java
       ctx.forward(
           record.withValue(
               KStreamsPayload.watermark(report.getWatermarkMillis(), 
transformId, 0, 1)));
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to