This is an automated email from the ASF dual-hosted git repository. tvalentyn pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/beam.git
commit 8a2ed4d7f56234a232a42485f6c83c42943b52c0 Merge: bc06581e110 36ba537c6c2 Author: tvalentyn <[email protected]> AuthorDate: Tue Oct 24 10:18:28 2023 -0700 Add MetricsContainer support to the Flink sources. #28609 Co-authored-by: Jiangjie Qin <[email protected]> .../streaming/io/source/SourceTestCompat.java | 13 ++ .../streaming/io/source/SourceTestCompat.java | 13 ++ .../beam/runners/flink/FlinkRunnerResult.java | 4 +- .../flink/metrics/FlinkMetricContainer.java | 160 +-------------------- ...ontainer.java => FlinkMetricContainerBase.java} | 74 +++------- .../FlinkMetricContainerWithoutAccumulator.java | 34 +++++ .../flink/metrics/ReaderInvocationUtil.java | 4 +- .../wrappers/streaming/io/source/FlinkSource.java | 12 +- .../streaming/io/source/FlinkSourceReaderBase.java | 13 +- .../io/source/bounded/FlinkBoundedSource.java | 8 +- .../source/bounded/FlinkBoundedSourceReader.java | 10 +- .../io/source/unbounded/FlinkUnboundedSource.java | 13 +- .../unbounded/FlinkUnboundedSourceReader.java | 6 +- .../flink/metrics/FlinkMetricContainerTest.java | 2 +- .../wrappers/streaming/io/TestCountingSource.java | 7 + .../io/source/FlinkSourceReaderTestBase.java | 28 ++++ .../bounded/FlinkBoundedSourceReaderTest.java | 5 +- .../unbounded/FlinkUnboundedSourceReaderTest.java | 5 +- 18 files changed, 176 insertions(+), 235 deletions(-)
