Port Flink integration tests to new DoFn
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/ae1f6d18 Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/ae1f6d18 Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/ae1f6d18 Branch: refs/heads/gearpump-runner Commit: ae1f6d181ebe3c0bdffc35c833a6fdc858937d6c Parents: 879f18f Author: Kenneth Knowles <[email protected]> Authored: Fri Aug 5 12:17:20 2016 -0700 Committer: Kenneth Knowles <[email protected]> Committed: Mon Aug 8 11:35:17 2016 -0700 ---------------------------------------------------------------------- .../java/org/apache/beam/runners/flink/ReadSourceITCase.java | 6 +++--- .../apache/beam/runners/flink/ReadSourceStreamingITCase.java | 8 +++++--- 2 files changed, 8 insertions(+), 6 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/ae1f6d18/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceITCase.java ---------------------------------------------------------------------- diff --git a/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceITCase.java b/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceITCase.java index ca70096..516c7ba 100644 --- a/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceITCase.java +++ b/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceITCase.java @@ -20,7 +20,7 @@ package org.apache.beam.runners.flink; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.CountingInput; import org.apache.beam.sdk.io.TextIO; -import org.apache.beam.sdk.transforms.OldDoFn; +import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.values.PCollection; @@ -72,8 +72,8 @@ public class ReadSourceITCase extends JavaProgramTestBase { PCollection<String> result = p .apply(CountingInput.upTo(10)) - .apply(ParDo.of(new OldDoFn<Long, String>() { - @Override + .apply(ParDo.of(new DoFn<Long, String>() { + @ProcessElement public void processElement(ProcessContext c) throws Exception { c.output(c.element().toString()); } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/ae1f6d18/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceStreamingITCase.java ---------------------------------------------------------------------- diff --git a/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceStreamingITCase.java b/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceStreamingITCase.java index bc69f34..ea58d0d 100644 --- a/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceStreamingITCase.java +++ b/runners/flink/runner/src/test/java/org/apache/beam/runners/flink/ReadSourceStreamingITCase.java @@ -20,9 +20,11 @@ package org.apache.beam.runners.flink; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.CountingInput; import org.apache.beam.sdk.io.TextIO; -import org.apache.beam.sdk.transforms.OldDoFn; +import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; + import com.google.common.base.Joiner; + import org.apache.flink.streaming.util.StreamingProgramTestBase; /** @@ -59,8 +61,8 @@ public class ReadSourceStreamingITCase extends StreamingProgramTestBase { p .apply(CountingInput.upTo(10)) - .apply(ParDo.of(new OldDoFn<Long, String>() { - @Override + .apply(ParDo.of(new DoFn<Long, String>() { + @ProcessElement public void processElement(ProcessContext c) throws Exception { c.output(c.element().toString()); }
