FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

Backport from Beam: Use UnboundedReadFromBoundedSource in dataflow ru… · drieber/DataflowJavaSDK@77b4b7b · GitHub

Repository navigation

Commit 77b4b7b

Browse files
authored andcommitted
Backport from Beam: Use UnboundedReadFromBoundedSource in dataflow runner (GoogleCloudPlatform#333)
1 parent ce2e5c0 commit 77b4b7b

4 files changed

Lines changed: 942 additions & 34 deletions

File tree

‎sdk/src/main/java/com/google/cloud/dataflow/sdk/runners/DataflowPipelineRunner.java‎

Lines changed: 31 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,7 @@
5757
import com.google.cloud.dataflow.sdk.coders.VarLongCoder;
5858
import com.google.cloud.dataflow.sdk.io.AvroIO;
5959
import com.google.cloud.dataflow.sdk.io.BigQueryIO;
60+
import com.google.cloud.dataflow.sdk.io.BoundedSource;
6061
import com.google.cloud.dataflow.sdk.io.FileBasedSink;
6162
import com.google.cloud.dataflow.sdk.io.PubsubIO;
6263
import com.google.cloud.dataflow.sdk.io.PubsubIO.Read.Bound.PubsubReader;
@@ -79,6 +80,7 @@
7980
import com.google.cloud.dataflow.sdk.runners.DataflowPipelineTranslator.TranslationContext;
8081
import com.google.cloud.dataflow.sdk.runners.dataflow.AssignWindows;
8182
import com.google.cloud.dataflow.sdk.runners.dataflow.DataflowAggregatorTransforms;
83+
import com.google.cloud.dataflow.sdk.runners.dataflow.DataflowUnboundedReadFromBoundedSource;
8284
import com.google.cloud.dataflow.sdk.runners.dataflow.ReadTranslator;
8385
import com.google.cloud.dataflow.sdk.runners.worker.IsmFormat;
8486
import com.google.cloud.dataflow.sdk.runners.worker.IsmFormat.IsmRecord;
@@ -354,11 +356,8 @@ public static DataflowPipelineRunner fromOptions(PipelineOptions options) {
354356
builder.put(View.AsIterable.class, StreamingViewAsIterable.class);
355357
builder.put(Write.Bound.class, StreamingWrite.class);
356358
builder.put(Read.Unbounded.class, StreamingUnboundedRead.class);
357-
builder.put(Read.Bounded.class, UnsupportedIO.class);
358-
builder.put(AvroIO.Read.Bound.class, UnsupportedIO.class);
359+
builder.put(Read.Bounded.class, StreamingBoundedRead.class);
359360
builder.put(AvroIO.Write.Bound.class, UnsupportedIO.class);
360-
builder.put(BigQueryIO.Read.Bound.class, UnsupportedIO.class);
361-
builder.put(TextIO.Read.Bound.class, UnsupportedIO.class);
362361
builder.put(TextIO.Write.Bound.class, UnsupportedIO.class);
363362
builder.put(Window.Bound.class, AssignWindows.class);
364363
// In streaming mode must use either the custom Pubsub unbounded source/sink or
@@ -2805,6 +2804,34 @@ public void processElement(ProcessContext c) {
28052804
}
28062805
}
28072806

2807+
/**
2808+
* Specialized implementation for
2809+
* {@link com.google.cloud.dataflow.sdk.io.Read.Bounded Read.Bounded} for the
2810+
* Dataflow runner in streaming mode.
2811+
*/
2812+
private static class StreamingBoundedRead<T> extends PTransform<PInput, PCollection<T>> {
2813+
private final BoundedSource<T> source;
2814+
2815+
/** Builds an instance of this class from the overridden transform. */
2816+
@SuppressWarnings("unused") // used via reflection in DataflowRunner#apply()
2817+
public StreamingBoundedRead(DataflowPipelineRunner runner, Read.Bounded<T> transform) {
2818+
this.source = transform.getSource();
2819+
}
2820+
2821+
@Override
2822+
protected Coder<T> getDefaultOutputCoder() {
2823+
return source.getDefaultOutputCoder();
2824+
}
2825+
2826+
@Override
2827+
public final PCollection<T> apply(PInput input) {
2828+
source.validate();
2829+
2830+
return Pipeline.applyTransform(input, new DataflowUnboundedReadFromBoundedSource<>(source))
2831+
.setIsBoundedInternal(IsBounded.BOUNDED);
2832+
}
2833+
}
2834+
28082835
/**
28092836
* Specialized implementation for
28102837
* {@link com.google.cloud.dataflow.sdk.transforms.Create.Values Create.Values} for the

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL