je-ik commented on code in PR #40186:
URL: https://github.com/apache/beam/pull/40186#discussion_r4060961028
##########
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:
Should this actually expand the target partitions and emit multiple elements
that target different partitions? Or is this done elsewhere?
--
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]