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]