【问题标题】:Throughput and Latency on Apache FlinkApache Flink 的吞吐量和延迟
【发布时间】:2017-11-19 03:16:37
【问题描述】:

我为 Apache Flink 编写了一个非常简单的 Java 程序,现在我对测量统计数据感兴趣,例如吞吐量(每秒处理的元组数)和延迟(程序处理每个输入元组所需的时间)。

 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

env.readTextFile("/home/LizardKing/Documents/Power/Prova.csv")
        .map(new MyMapper().writeAsCsv("/home/LizardKing/Results.csv");

JobExecutionResult res = env.execute();

我知道 Flink 公开了一些指标:

https://ci.apache.org/projects/flink/flink-docs-release-1.2/monitoring/metrics.html

但我不确定如何使用它们来获得我想要的东西。从链接中我读到可以使用“仪表”来测量平均吞吐量,但是在定义它之后,我应该如何使用它?

【问题讨论】:

  • 你到底在纠结什么?对于吞吐量,您可以在 MyMapper 函数中注册一个 Meter,正如您提供的链接中所示。您可以在 Flink Web 仪表板中实时查看指标。
  • 如果我按照我需要实现 myMeter 类的说明进行操作,我已经尝试了一些方法,但它不起作用。如果我使用 DropWizard 仪表并尝试在独立模式下运行它,即使我在 pom.xml 中包含了依赖项,也会出现错误(java.lang.NoClassDefFoundError: com/codahale/metrics/Meter)。

标签: java apache-flink latency flink-streaming throughput


【解决方案1】:

我们正在运行在 yarn 上的生产流作业中运行自定义指标,例如 Meter、Gauge。

以下是步骤:

对 pom.xml 的附加依赖

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-metrics-dropwizard</artifactId>
    <version>${flink.version}</version>
</dependency>

我们使用的是 1.2.1 版

然后将meter添加到MyMapper类。

import org.apache.flink.api.common.JobExecutionResult;
import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.dropwizard.metrics.DropwizardMeterWrapper;
import org.apache.flink.metrics.Meter;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;


public class Test {


    public static void main(String[] args) throws Exception {

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env
                .readTextFile("/home/LizardKing/Documents/Power/Prova.csv")
                .map(new MyMapper())
                .writeAsCsv("/home/LizardKing/Results.csv");

        JobExecutionResult res = env.execute();
    }


    private static class MyMapper extends RichMapFunction<String, Object> {

        private transient Meter meter;

        @Override
        public void open(Configuration parameters) throws Exception {
            super.open(parameters);
            this.meter = getRuntimeContext()
                    .getMetricGroup()
                    .meter("myMeter", new DropwizardMeterWrapper(new com.codahale.metrics.Meter()));
        }

        @Override
        public Object map(String value) throws Exception {    
            this.meter.markEvent();
            return value;
        }
    }
}

希望这会有所帮助。

【讨论】:

  • 它有帮助,我还遇到了另一个问题:当我尝试在 flink(而不是从 IDE)中运行这个程序时,我发现在 pom.xml 中包含依赖项是不够的。我必须提供库来 flink,我被建议的方法是使用 maven-shade 插件。它应该将依赖项打包到上传的 jar 中。
猜你喜欢
  • 2017-10-29
  • 2022-12-25
  • 2018-12-14
  • 2015-04-16
  • 2021-04-16
  • 2019-08-27
  • 2016-11-16
  • 2017-03-05
  • 1970-01-01
相关资源
最近更新 更多