This closes #2576
Project: http://git-wip-us.apache.org/repos/asf/beam/repo Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/83193698 Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/83193698 Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/83193698 Branch: refs/heads/DSL_SQL Commit: 83193698d8ea3dc9cb2a3ed8fe6b4ee6b810237c Parents: 8a00f22 cdd2544 Author: Ismaël MejÃa <ieme...@apache.org> Authored: Wed Apr 19 15:07:54 2017 +0200 Committer: Ismaël MejÃa <ieme...@apache.org> Committed: Wed Apr 19 15:07:54 2017 +0200 ---------------------------------------------------------------------- ...PostCommit_Java_ValidatesRunner_Flink.groovy | 2 +- runners/flink/examples/pom.xml | 130 --- .../beam/runners/flink/examples/TFIDF.java | 455 -------- .../beam/runners/flink/examples/WordCount.java | 129 --- .../runners/flink/examples/package-info.java | 22 - .../flink/examples/streaming/AutoComplete.java | 400 ------- .../flink/examples/streaming/JoinExamples.java | 154 --- .../examples/streaming/WindowedWordCount.java | 141 --- .../flink/examples/streaming/package-info.java | 22 - runners/flink/pom.xml | 275 ++++- runners/flink/runner/pom.xml | 330 ------ .../flink/DefaultParallelismFactory.java | 39 - .../flink/FlinkBatchPipelineTranslator.java | 139 --- .../flink/FlinkBatchTransformTranslators.java | 723 ------------ .../flink/FlinkBatchTranslationContext.java | 153 --- .../flink/FlinkDetachedRunnerResult.java | 75 -- .../FlinkPipelineExecutionEnvironment.java | 241 ---- .../runners/flink/FlinkPipelineOptions.java | 101 -- .../runners/flink/FlinkPipelineTranslator.java | 53 - .../apache/beam/runners/flink/FlinkRunner.java | 232 ---- .../runners/flink/FlinkRunnerRegistrar.java | 62 -- .../beam/runners/flink/FlinkRunnerResult.java | 98 -- .../flink/FlinkStreamingPipelineTranslator.java | 276 ----- .../FlinkStreamingTransformTranslators.java | 1044 ----------------- .../flink/FlinkStreamingTranslationContext.java | 130 --- .../flink/FlinkStreamingViewOverrides.java | 372 ------- .../flink/PipelineTranslationOptimizer.java | 72 -- .../beam/runners/flink/TestFlinkRunner.java | 84 -- .../beam/runners/flink/TranslationMode.java | 31 - .../apache/beam/runners/flink/package-info.java | 22 - .../functions/FlinkAggregatorFactory.java | 53 - .../functions/FlinkAssignContext.java | 63 -- .../functions/FlinkAssignWindows.java | 49 - .../functions/FlinkDoFnFunction.java | 161 --- .../FlinkMergingNonShuffleReduceFunction.java | 228 ---- .../FlinkMergingPartialReduceFunction.java | 201 ---- .../functions/FlinkMergingReduceFunction.java | 199 ---- .../FlinkMultiOutputPruningFunction.java | 50 - .../functions/FlinkNoOpStepContext.java | 73 -- .../functions/FlinkPartialReduceFunction.java | 172 --- .../functions/FlinkReduceFunction.java | 173 --- .../functions/FlinkSideInputReader.java | 80 -- .../functions/FlinkStatefulDoFnFunction.java | 198 ---- .../functions/SideInputInitializer.java | 73 -- .../translation/functions/package-info.java | 22 - .../runners/flink/translation/package-info.java | 22 - .../translation/types/CoderTypeInformation.java | 120 -- .../translation/types/CoderTypeSerializer.java | 132 --- .../types/EncodedValueComparator.java | 195 ---- .../types/EncodedValueSerializer.java | 113 -- .../types/EncodedValueTypeInformation.java | 98 -- .../types/InspectableByteArrayOutputStream.java | 34 - .../flink/translation/types/KvKeySelector.java | 50 - .../flink/translation/types/package-info.java | 22 - .../utils/SerializedPipelineOptions.java | 67 -- .../flink/translation/utils/package-info.java | 22 - .../wrappers/DataInputViewWrapper.java | 58 - .../wrappers/DataOutputViewWrapper.java | 51 - .../SerializableFnAggregatorWrapper.java | 98 -- .../translation/wrappers/SourceInputFormat.java | 150 --- .../translation/wrappers/SourceInputSplit.java | 52 - .../translation/wrappers/package-info.java | 22 - .../wrappers/streaming/DoFnOperator.java | 774 ------------- .../streaming/KvToByteBufferKeySelector.java | 56 - .../streaming/SingletonKeyedWorkItem.java | 56 - .../streaming/SingletonKeyedWorkItemCoder.java | 126 --- .../streaming/SplittableDoFnOperator.java | 150 --- .../wrappers/streaming/WindowDoFnOperator.java | 117 -- .../wrappers/streaming/WorkItemKeySelector.java | 56 - .../streaming/io/BoundedSourceWrapper.java | 218 ---- .../streaming/io/UnboundedSocketSource.java | 249 ----- .../streaming/io/UnboundedSourceWrapper.java | 476 -------- .../wrappers/streaming/io/package-info.java | 22 - .../wrappers/streaming/package-info.java | 22 - .../state/FlinkBroadcastStateInternals.java | 865 -------------- .../state/FlinkKeyGroupStateInternals.java | 487 -------- .../state/FlinkSplitStateInternals.java | 260 ----- .../streaming/state/FlinkStateInternals.java | 1053 ------------------ .../state/KeyGroupCheckpointedOperator.java | 35 - .../state/KeyGroupRestoringOperator.java | 32 - .../wrappers/streaming/state/package-info.java | 22 - .../runner/src/main/resources/log4j.properties | 23 - .../flink/EncodedValueComparatorTest.java | 70 -- .../runners/flink/FlinkRunnerRegistrarTest.java | 48 - .../beam/runners/flink/FlinkTestPipeline.java | 72 -- .../beam/runners/flink/PipelineOptionsTest.java | 184 --- .../beam/runners/flink/ReadSourceITCase.java | 85 -- .../flink/ReadSourceStreamingITCase.java | 74 -- .../beam/runners/flink/WriteSinkITCase.java | 192 ---- .../flink/streaming/DoFnOperatorTest.java | 600 ---------- .../FlinkBroadcastStateInternalsTest.java | 245 ---- .../FlinkKeyGroupStateInternalsTest.java | 262 ----- .../streaming/FlinkSplitStateInternalsTest.java | 101 -- .../streaming/FlinkStateInternalsTest.java | 395 ------- .../flink/streaming/GroupByNullKeyTest.java | 124 --- .../flink/streaming/TestCountingSource.java | 254 ----- .../streaming/TopWikipediaSessionsITCase.java | 133 --- .../streaming/UnboundedSourceWrapperTest.java | 464 -------- .../runners/flink/streaming/package-info.java | 22 - .../src/test/resources/log4j-test.properties | 27 - .../flink/DefaultParallelismFactory.java | 39 + .../flink/FlinkBatchPipelineTranslator.java | 139 +++ .../flink/FlinkBatchTransformTranslators.java | 723 ++++++++++++ .../flink/FlinkBatchTranslationContext.java | 153 +++ .../flink/FlinkDetachedRunnerResult.java | 75 ++ .../FlinkPipelineExecutionEnvironment.java | 241 ++++ .../runners/flink/FlinkPipelineOptions.java | 101 ++ .../runners/flink/FlinkPipelineTranslator.java | 53 + .../apache/beam/runners/flink/FlinkRunner.java | 232 ++++ .../runners/flink/FlinkRunnerRegistrar.java | 62 ++ .../beam/runners/flink/FlinkRunnerResult.java | 98 ++ .../flink/FlinkStreamingPipelineTranslator.java | 276 +++++ .../FlinkStreamingTransformTranslators.java | 1044 +++++++++++++++++ .../flink/FlinkStreamingTranslationContext.java | 130 +++ .../flink/FlinkStreamingViewOverrides.java | 372 +++++++ .../flink/PipelineTranslationOptimizer.java | 72 ++ .../beam/runners/flink/TestFlinkRunner.java | 84 ++ .../beam/runners/flink/TranslationMode.java | 31 + .../apache/beam/runners/flink/package-info.java | 22 + .../functions/FlinkAggregatorFactory.java | 53 + .../functions/FlinkAssignContext.java | 63 ++ .../functions/FlinkAssignWindows.java | 49 + .../functions/FlinkDoFnFunction.java | 161 +++ .../FlinkMergingNonShuffleReduceFunction.java | 228 ++++ .../FlinkMergingPartialReduceFunction.java | 201 ++++ .../functions/FlinkMergingReduceFunction.java | 199 ++++ .../FlinkMultiOutputPruningFunction.java | 50 + .../functions/FlinkNoOpStepContext.java | 73 ++ .../functions/FlinkPartialReduceFunction.java | 172 +++ .../functions/FlinkReduceFunction.java | 173 +++ .../functions/FlinkSideInputReader.java | 80 ++ .../functions/FlinkStatefulDoFnFunction.java | 198 ++++ .../functions/SideInputInitializer.java | 73 ++ .../translation/functions/package-info.java | 22 + .../runners/flink/translation/package-info.java | 22 + .../translation/types/CoderTypeInformation.java | 120 ++ .../translation/types/CoderTypeSerializer.java | 132 +++ .../types/EncodedValueComparator.java | 195 ++++ .../types/EncodedValueSerializer.java | 113 ++ .../types/EncodedValueTypeInformation.java | 98 ++ .../types/InspectableByteArrayOutputStream.java | 34 + .../flink/translation/types/KvKeySelector.java | 50 + .../flink/translation/types/package-info.java | 22 + .../utils/SerializedPipelineOptions.java | 67 ++ .../flink/translation/utils/package-info.java | 22 + .../wrappers/DataInputViewWrapper.java | 58 + .../wrappers/DataOutputViewWrapper.java | 51 + .../SerializableFnAggregatorWrapper.java | 98 ++ .../translation/wrappers/SourceInputFormat.java | 150 +++ .../translation/wrappers/SourceInputSplit.java | 52 + .../translation/wrappers/package-info.java | 22 + .../wrappers/streaming/DoFnOperator.java | 774 +++++++++++++ .../streaming/KvToByteBufferKeySelector.java | 56 + .../streaming/SingletonKeyedWorkItem.java | 56 + .../streaming/SingletonKeyedWorkItemCoder.java | 126 +++ .../streaming/SplittableDoFnOperator.java | 150 +++ .../wrappers/streaming/WindowDoFnOperator.java | 117 ++ .../wrappers/streaming/WorkItemKeySelector.java | 56 + .../streaming/io/BoundedSourceWrapper.java | 218 ++++ .../streaming/io/UnboundedSocketSource.java | 249 +++++ .../streaming/io/UnboundedSourceWrapper.java | 476 ++++++++ .../wrappers/streaming/io/package-info.java | 22 + .../wrappers/streaming/package-info.java | 22 + .../state/FlinkBroadcastStateInternals.java | 865 ++++++++++++++ .../state/FlinkKeyGroupStateInternals.java | 487 ++++++++ .../state/FlinkSplitStateInternals.java | 260 +++++ .../streaming/state/FlinkStateInternals.java | 1053 ++++++++++++++++++ .../state/KeyGroupCheckpointedOperator.java | 35 + .../state/KeyGroupRestoringOperator.java | 32 + .../wrappers/streaming/state/package-info.java | 22 + .../flink/src/main/resources/log4j.properties | 23 + .../flink/EncodedValueComparatorTest.java | 70 ++ .../runners/flink/FlinkRunnerRegistrarTest.java | 48 + .../beam/runners/flink/FlinkTestPipeline.java | 72 ++ .../beam/runners/flink/PipelineOptionsTest.java | 184 +++ .../beam/runners/flink/ReadSourceITCase.java | 85 ++ .../flink/ReadSourceStreamingITCase.java | 74 ++ .../beam/runners/flink/WriteSinkITCase.java | 192 ++++ .../flink/streaming/DoFnOperatorTest.java | 600 ++++++++++ .../FlinkBroadcastStateInternalsTest.java | 245 ++++ .../FlinkKeyGroupStateInternalsTest.java | 262 +++++ .../streaming/FlinkSplitStateInternalsTest.java | 101 ++ .../streaming/FlinkStateInternalsTest.java | 395 +++++++ .../flink/streaming/GroupByNullKeyTest.java | 124 +++ .../flink/streaming/TestCountingSource.java | 254 +++++ .../streaming/TopWikipediaSessionsITCase.java | 133 +++ .../streaming/UnboundedSourceWrapperTest.java | 464 ++++++++ .../runners/flink/streaming/package-info.java | 22 + .../src/test/resources/log4j-test.properties | 27 + 189 files changed, 15765 insertions(+), 17293 deletions(-) ----------------------------------------------------------------------