junaiddshaukat commented on code in PR #40186:
URL: https://github.com/apache/beam/pull/40186#discussion_r4063377910


##########
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ShuffleByKeyProcessor.java:
##########
@@ -114,6 +123,15 @@ public void process(Record<byte[], KStreamsPayload<?>> 
record) {
         throw new RuntimeException("Failed to encode shuffle key", e);
       }
       ctx.forward(record.withKey(encodedKey));
+    } else if (payload.isFlush()) {
+      // Retargeted for this edge. Fanning in, an instance may have nothing to 
address.
+      Set<Integer> targets =
+          flushTargets(upstreamPartition, upstreamPartitionCount, 
downstreamPartitionCount);
+      if (!targets.isEmpty()) {
+        ctx.forward(

Review Comment:
   It is done in the sink, by Kafka Streams. The shuffle forwards the flush 
once, `KStreamsPayloadPartitioner.partitions()` returns the target set, and 
`RecordCollectorImpl` sends one copy to each partition in it 
([RecordCollectorImpl.java#L163-L177](https://github.com/apache/kafka/blob/3.9.0/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordCollectorImpl.java#L163-L177),
 3.9.0):
   
   ```java
   for (final int multicastPartition: multicastPartitions) {
       send(topic, key, value, headers, multicastPartition, timestamp, ...);
   }
   ```
   
   This is the multicast from KIP-837, and it is the same path the watermark 
broadcast already uses, only with a subset instead of every partition. An empty 
set is dropped there with a warning, which is why the shuffle does not forward 
one.
   
   Expanding it here would also work, but it would do the same thing with more 
records, since a processor cannot pick the partition itself; only the sink's 
partitioner can.
   
   It was not obvious from the code, so I added a short comment here, and a 
broker test (`KStreamsPayloadPartitionerBrokerIT`) that sends a flush for 
partitions 1 and 2 of a 4 partition topic and checks it arrives once on each 
and nowhere else. The unit tests cannot show this, because `TopologyTestDriver` 
reports every topic as having one partition.
   



-- 
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