This closes #2300
Project: http://git-wip-us.apache.org/repos/asf/beam/repo Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/8160924e Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/8160924e Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/8160924e Branch: refs/heads/master Commit: 8160924e12a4723df34209a04c907d4f95466ecf Parents: b382795 217f24b Author: Aljoscha Krettek <[email protected]> Authored: Fri Apr 21 11:59:46 2017 +0200 Committer: Aljoscha Krettek <[email protected]> Committed: Fri Apr 21 11:59:46 2017 +0200 ---------------------------------------------------------------------- runners/flink/pom.xml | 8 +- .../flink/FlinkBatchTransformTranslators.java | 7 +- .../flink/FlinkBatchTranslationContext.java | 4 + .../runners/flink/FlinkPipelineOptions.java | 5 + .../beam/runners/flink/FlinkRunnerResult.java | 3 +- .../flink/FlinkStreamingPipelineTranslator.java | 3 + .../FlinkStreamingTransformTranslators.java | 15 + .../flink/FlinkStreamingTranslationContext.java | 3 + .../metrics/DoFnRunnerWithMetricsUpdate.java | 91 ++++++ .../flink/metrics/FlinkMetricContainer.java | 315 +++++++++++++++++++ .../flink/metrics/FlinkMetricResults.java | 146 +++++++++ .../flink/metrics/ReaderInvocationUtil.java | 71 +++++ .../runners/flink/metrics/package-info.java | 22 ++ .../functions/FlinkDoFnFunction.java | 10 + .../functions/FlinkStatefulDoFnFunction.java | 10 + .../translation/wrappers/SourceInputFormat.java | 20 +- .../wrappers/streaming/DoFnOperator.java | 13 +- .../streaming/SplittableDoFnOperator.java | 2 + .../wrappers/streaming/WindowDoFnOperator.java | 2 + .../streaming/io/BoundedSourceWrapper.java | 17 +- .../streaming/io/UnboundedSourceWrapper.java | 18 +- .../beam/runners/flink/PipelineOptionsTest.java | 10 + .../flink/streaming/DoFnOperatorTest.java | 5 + .../streaming/UnboundedSourceWrapperTest.java | 12 +- 24 files changed, 788 insertions(+), 24 deletions(-) ----------------------------------------------------------------------
