【发布时间】:2020-07-25 09:10:35
【问题描述】:
随着越来越多的记录被处理,我的程序变得非常慢。我最初认为这是由于内存消耗过多,因为我的程序是字符串密集型的(我使用的是 Java 11,所以应该尽可能使用紧凑的字符串)所以我增加了 JVM 堆:
-Xms2048m
-Xmx6144m
我还增加了任务管理器的内存以及超时,flink-conf.yaml:
jobmanager.heap.size: 6144m
heartbeat.timeout: 5000000
但是,这些都没有帮助解决这个问题。在处理大约 350 万条记录之后,该程序仍然变得非常缓慢,只剩下大约 50 万条记录。随着程序接近 350 万大关,它变得非常非常慢,直到最终超时,总执行时间约为 11 分钟。
我在 VisualVm 中检查了内存消耗,但内存消耗从未超过 700MB。我的 flink 管道如下所示:
final StreamExecutionEnvironment environment = StreamExecutionEnvironment.createLocalEnvironment(1);
environment.setParallelism(1);
DataStream<Tuple> stream = environment.addSource(new TPCHQuery3Source(filePaths, relations));
stream.process(new TPCHQuery3Process(relations)).addSink(new FDSSink());
environment.execute("FlinkDataService");
大部分工作是在 process 函数中完成的,我正在实现数据库连接算法,并且列存储为字符串,具体来说,我正在实现 TPCH 基准的查询 3,如果您愿意,请点击此处https://examples.citusdata.com/tpch_queries.html .
超时错误是这样的:
java.util.concurrent.TimeoutException: Heartbeat of TaskManager with id <id> timed out.
一旦我也收到此错误:
Exception in thread "pool-1-thread-1" java.lang.OutOfMemoryError: Java heap space
另外,我的 VisualVM 监控屏幕截图是在事情变得非常缓慢的时候捕获的:
这是我的源函数的运行循环:
while (run) {
readers.forEach(reader -> {
try {
String line = reader.readLine();
if (line != null) {
Tuple tuple = lineToTuple(line, counter.get() % filePaths.size());
if (tuple != null && isValidTuple(tuple)) {
sourceContext.collect(tuple);
}
} else {
closedReaders.add(reader);
if (closedReaders.size() == filePaths.size()) {
System.out.println("ALL FILES HAVE BEEN STREAMED");
cancel();
}
}
counter.getAndIncrement();
} catch (IOException e) {
e.printStackTrace();
}
});
}
我基本上读取了我需要的 3 个文件中的每一行,根据文件的顺序,我构造了一个元组对象,它是我的自定义类,称为元组,表示表中的一行,如果它发出该元组有效,即满足日期的某些条件。
我还建议 JVM 在第 100 万、150 万、200 万和 250 万记录处进行垃圾收集,如下所示:
System.gc()
有什么想法可以优化这个吗?
【问题讨论】:
-
我使用 Flink 的数据流 API 实现了 TPC-H 查询 03,我只流式传输 Order 数据集。我用作状态的 Customer 和 Lineitem 数据集。通过这样做,我只需要增加超时和 JVM 内存。 github.com/felipegutierrez/explore-flink/blob/master/src/main/…
-
你的程序在你的 IDE 和 flink 集群上运行是否同样快?
-
不,当然不是。在 IDE 中,我必须减小 LineItem 表的大小,因为它太大了 ~725MB。尽管如此,我还是可以从 IDE 运行。
-
嗯,我的程序在集群上很慢,不知道为什么..
-
我觉得我需要把 JVM 的堆大小调大,你知道 flink 集群怎么做吗?
标签: java timeout apache-flink taskmanager