【发布时间】: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