This is an automated email from the ASF dual-hosted git repository.
je-ik pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
from 7d0f1bf3870 Fix PickleCoder.as_deterministic_coder() raising TypeError
(#39943)
add cf2c5767e32 Kafka Streams runner: target flush markers across
repartition topics
add e038146a76b Kafka Streams runner: check flush multicast on a real
broker
new 25dc1ba1da8 Merge pull request #40186: Kafka Streams runner: target
flush markers across repartition topics
The 1 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.../streams/translation/GroupByKeyTranslator.java | 14 +-
...tioner.java => KStreamsPayloadPartitioner.java} | 25 ++-
.../streams/translation/ShuffleByKeyProcessor.java | 35 +++-
.../streams/translation/TerminationTracker.java | 4 +-
.../KStreamsPayloadPartitionerBrokerIT.java | 176 +++++++++++++++++++++
.../KStreamsPayloadPartitionerTest.java | 74 +++++++++
.../translation/ShuffleByKeyProcessorTest.java | 75 ++++++++-
7 files changed, 389 insertions(+), 14 deletions(-)
rename
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/{GroupByKeyBroadcastPartitioner.java
=> KStreamsPayloadPartitioner.java} (74%)
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadPartitionerBrokerIT.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadPartitionerTest.java