| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent ce2e5c0 commit 77b4b7b
4 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -57,6 +57,7 @@ | |||
| 57 | 57 | import com.google.cloud.dataflow.sdk.coders.VarLongCoder; | |
| 58 | 58 | import com.google.cloud.dataflow.sdk.io.AvroIO; | |
| 59 | 59 | import com.google.cloud.dataflow.sdk.io.BigQueryIO; | |
| 60 | + import com.google.cloud.dataflow.sdk.io.BoundedSource; | ||
| 60 | 61 | import com.google.cloud.dataflow.sdk.io.FileBasedSink; | |
| 61 | 62 | import com.google.cloud.dataflow.sdk.io.PubsubIO; | |
| 62 | 63 | import com.google.cloud.dataflow.sdk.io.PubsubIO.Read.Bound.PubsubReader; | |
@@ -79,6 +80,7 @@ | |||
| 79 | 80 | import com.google.cloud.dataflow.sdk.runners.DataflowPipelineTranslator.TranslationContext; | |
| 80 | 81 | import com.google.cloud.dataflow.sdk.runners.dataflow.AssignWindows; | |
| 81 | 82 | import com.google.cloud.dataflow.sdk.runners.dataflow.DataflowAggregatorTransforms; | |
| 83 | + import com.google.cloud.dataflow.sdk.runners.dataflow.DataflowUnboundedReadFromBoundedSource; | ||
| 82 | 84 | import com.google.cloud.dataflow.sdk.runners.dataflow.ReadTranslator; | |
| 83 | 85 | import com.google.cloud.dataflow.sdk.runners.worker.IsmFormat; | |
| 84 | 86 | import com.google.cloud.dataflow.sdk.runners.worker.IsmFormat.IsmRecord; | |
@@ -354,11 +356,8 @@ public static DataflowPipelineRunner fromOptions(PipelineOptions options) { | |||
| 354 | 356 | builder.put(View.AsIterable.class, StreamingViewAsIterable.class); | |
| 355 | 357 | builder.put(Write.Bound.class, StreamingWrite.class); | |
| 356 | 358 | 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); | ||
| 359 | 360 | 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); | ||
| 362 | 361 | builder.put(TextIO.Write.Bound.class, UnsupportedIO.class); | |
| 363 | 362 | builder.put(Window.Bound.class, AssignWindows.class); | |
| 364 | 363 | // In streaming mode must use either the custom Pubsub unbounded source/sink or | |
@@ -2805,6 +2804,34 @@ public void processElement(ProcessContext c) { | |||
| 2805 | 2804 | } | |
| 2806 | 2805 | } | |
| 2807 | 2806 | ||
| 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 | + | ||
| 2808 | 2835 | /** | |
| 2809 | 2836 | * Specialized implementation for | |
| 2810 | 2837 | * {@link com.google.cloud.dataflow.sdk.transforms.Create.Values Create.Values} for the | |
| Back | FazBrowse Home | New Git URL |
0 commit comments