This closes #737
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/603f337b Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/603f337b Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/603f337b Branch: refs/heads/master Commit: 603f337b130cf9bb3a8e2a810db99ed62211c2ef Parents: a87015b 695a80a Author: Kenneth Knowles <[email protected]> Authored: Wed Aug 24 12:47:54 2016 -0700 Committer: Kenneth Knowles <[email protected]> Committed: Wed Aug 24 12:47:54 2016 -0700 ---------------------------------------------------------------------- .../beam/runners/core/SideInputHandler.java | 240 ++++ .../beam/runners/core/SideInputHandlerTest.java | 222 ++++ runners/flink/runner/pom.xml | 11 +- .../apache/beam/runners/flink/FlinkRunner.java | 386 ++++++- .../beam/runners/flink/TestFlinkRunner.java | 23 +- .../FlinkStreamingPipelineTranslator.java | 59 +- .../FlinkStreamingTransformTranslators.java | 952 ++++++++++++---- .../functions/FlinkDoFnFunction.java | 4 +- .../functions/FlinkMultiOutputDoFnFunction.java | 7 +- .../translation/types/CoderTypeInformation.java | 4 + .../wrappers/streaming/DoFnOperator.java | 516 +++++++++ .../streaming/FlinkAbstractParDoWrapper.java | 282 ----- .../FlinkGroupAlsoByWindowWrapper.java | 644 ----------- .../streaming/FlinkGroupByKeyWrapper.java | 73 -- .../streaming/FlinkParDoBoundMultiWrapper.java | 79 -- .../streaming/FlinkParDoBoundWrapper.java | 104 -- .../wrappers/streaming/FlinkStateInternals.java | 1038 ++++++++++++++++++ .../streaming/SingletonKeyedWorkItem.java | 54 + .../streaming/SingletonKeyedWorkItemCoder.java | 125 +++ .../wrappers/streaming/WindowDoFnOperator.java | 345 ++++++ .../wrappers/streaming/WorkItemKeySelector.java | 58 + .../streaming/io/BoundedSourceWrapper.java | 219 ++++ .../io/FlinkStreamingCreateFunction.java | 56 - .../state/AbstractFlinkTimerInternals.java | 127 --- .../streaming/state/FlinkStateInternals.java | 733 ------------- .../streaming/state/StateCheckpointReader.java | 93 -- .../streaming/state/StateCheckpointUtils.java | 155 --- .../streaming/state/StateCheckpointWriter.java | 131 --- .../wrappers/streaming/state/StateType.java | 73 -- .../beam/runners/flink/PipelineOptionsTest.java | 103 +- .../flink/streaming/DoFnOperatorTest.java | 328 ++++++ .../streaming/FlinkStateInternalsTest.java | 391 +++++++ .../flink/streaming/GroupAlsoByWindowTest.java | 523 --------- .../flink/streaming/StateSerializationTest.java | 338 ------ .../streaming/UnboundedSourceWrapperTest.java | 2 +- .../org/apache/beam/sdk/testing/PAssert.java | 22 +- .../beam/sdk/transforms/join/RawUnionValue.java | 25 + .../apache/beam/sdk/transforms/CombineTest.java | 11 +- .../beam/sdk/transforms/ParDoLifecycleTest.java | 3 +- .../apache/beam/sdk/transforms/ParDoTest.java | 4 +- 40 files changed, 4810 insertions(+), 3753 deletions(-) ----------------------------------------------------------------------
