Skip to content

Commit

Permalink
Revert "Add Redistribute translation to Samza runner"
Browse files Browse the repository at this point in the history
This reverts commit 8f1d3da.
  • Loading branch information
kennknowles committed Apr 30, 2024
1 parent ece456c commit ab9af72
Show file tree
Hide file tree
Showing 3 changed files with 2 additions and 77 deletions.

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -42,16 +42,6 @@
public class ReshuffleTranslator<K, InT, OutT>
implements TransformTranslator<PTransform<PCollection<KV<K, InT>>, PCollection<KV<K, OutT>>>> {

private final String prefix;

ReshuffleTranslator(String prefix) {
this.prefix = prefix;
}

ReshuffleTranslator() {
this("rshfl-");
}

@Override
public void translate(
PTransform<PCollection<KV<K, InT>>, PCollection<KV<K, OutT>>> transform,
Expand All @@ -70,7 +60,7 @@ public void translate(
inputStream,
inputCoder.getKeyCoder(),
elementCoder,
prefix + ctx.getTransformId(),
"rshfl-" + ctx.getTransformId(),
ctx.getPipelineOptions().getMaxSourceParallelism() > 1);

ctx.registerMessageStream(output, outputStream);
Expand All @@ -93,7 +83,7 @@ public void translatePortable(
inputStream,
((KvCoder<K, InT>) windowedInputCoder.getValueCoder()).getKeyCoder(),
windowedInputCoder,
prefix + ctx.getTransformId(),
"rshfl-" + ctx.getTransformId(),
ctx.getPipelineOptions().getMaxSourceParallelism() > 1);

ctx.registerMessageStream(outputId, outputStream);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -179,7 +179,6 @@ public Map<String, TransformTranslator<?>> getTransformTranslators() {
return ImmutableMap.<String, TransformTranslator<?>>builder()
.put(PTransformTranslation.READ_TRANSFORM_URN, new ReadTranslator<>())
.put(PTransformTranslation.RESHUFFLE_URN, new ReshuffleTranslator<>())
.put(PTransformTranslation.REDISTRIBUTE_BY_KEY_URN, new RedistributeByKeyTranslator<>())
.put(PTransformTranslation.PAR_DO_TRANSFORM_URN, new ParDoBoundMultiTranslator<>())
.put(PTransformTranslation.GROUP_BY_KEY_TRANSFORM_URN, new GroupByKeyTranslator<>())
.put(PTransformTranslation.COMBINE_PER_KEY_TRANSFORM_URN, new GroupByKeyTranslator<>())
Expand Down

0 comments on commit ab9af72

Please sign in to comment.