【问题标题】:Flink Task Manager timeoutFlink 任务管理器超时
【发布时间】: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


【解决方案1】:

字符串intern()救了我。在将每个字符串存储到我的地图之前,我对每个字符串进行了实习,这就像一个魅力。

【讨论】:

    【解决方案2】:

    这些是我在链接独立集群上更改的属性,用于计算 TPC-H 查询 03。

    jobmanager.memory.process.size: 1600m
    heartbeat.timeout: 100000
    taskmanager.memory.process.size: 8g # defaul: 1728m
    

    我实现了这个查询以仅对 Order 表进行流式传输,并将其他表保留为状态。此外,我将计算作为无窗口查询,我认为它更有意义并且速度更快。

    public class TPCHQuery03 {
    
        private final String topic = "topic-tpch-query-03";
    
        public TPCHQuery03() {
            this(PARAMETER_OUTPUT_LOG, "127.0.0.1", false, false, -1);
        }
    
        public TPCHQuery03(String output, String ipAddressSink, boolean disableOperatorChaining, boolean pinningPolicy, long maxCount) {
            try {
                StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
                env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime);
    
                if (disableOperatorChaining) {
                    env.disableOperatorChaining();
                }
    
                DataStream<Order> orders = env
                        .addSource(new OrdersSource(maxCount)).name(OrdersSource.class.getSimpleName()).uid(OrdersSource.class.getSimpleName());
    
                // Filter market segment "AUTOMOBILE"
                // customers = customers.filter(new CustomerFilter());
    
                // Filter all Orders with o_orderdate < 12.03.1995
                DataStream<Order> ordersFiltered = orders
                        .filter(new OrderDateFilter("1995-03-12")).name(OrderDateFilter.class.getSimpleName()).uid(OrderDateFilter.class.getSimpleName());
    
                // Join customers with orders and package them into a ShippingPriorityItem
                DataStream<ShippingPriorityItem> customerWithOrders = ordersFiltered
                        .keyBy(new OrderKeySelector())
                        .process(new OrderKeyedByCustomerProcessFunction(pinningPolicy)).name(OrderKeyedByCustomerProcessFunction.class.getSimpleName()).uid(OrderKeyedByCustomerProcessFunction.class.getSimpleName());
    
                // Join the last join result with Lineitems
                DataStream<ShippingPriorityItem> result = customerWithOrders
                        .keyBy(new ShippingPriorityOrderKeySelector())
                        .process(new ShippingPriorityKeyedProcessFunction(pinningPolicy)).name(ShippingPriorityKeyedProcessFunction.class.getSimpleName()).uid(ShippingPriorityKeyedProcessFunction.class.getSimpleName());
    
                // Group by l_orderkey, o_orderdate and o_shippriority and compute revenue sum
                DataStream<ShippingPriorityItem> resultSum = result
                        .keyBy(new ShippingPriority3KeySelector())
                        .reduce(new SumShippingPriorityItem(pinningPolicy)).name(SumShippingPriorityItem.class.getSimpleName()).uid(SumShippingPriorityItem.class.getSimpleName());
    
                // emit result
                if (output.equalsIgnoreCase(PARAMETER_OUTPUT_MQTT)) {
                    resultSum
                            .map(new ShippingPriorityItemMap(pinningPolicy)).name(ShippingPriorityItemMap.class.getSimpleName()).uid(ShippingPriorityItemMap.class.getSimpleName())
                            .addSink(new MqttStringPublisher(ipAddressSink, topic, pinningPolicy)).name(OPERATOR_SINK).uid(OPERATOR_SINK);
                } else if (output.equalsIgnoreCase(PARAMETER_OUTPUT_LOG)) {
                    resultSum.print().name(OPERATOR_SINK).uid(OPERATOR_SINK);
                } else if (output.equalsIgnoreCase(PARAMETER_OUTPUT_FILE)) {
                    StreamingFileSink<String> sink = StreamingFileSink
                            .forRowFormat(new Path(PATH_OUTPUT_FILE), new SimpleStringEncoder<String>("UTF-8"))
                            .withRollingPolicy(
                                    DefaultRollingPolicy.builder().withRolloverInterval(TimeUnit.MINUTES.toMillis(15))
                                            .withInactivityInterval(TimeUnit.MINUTES.toMillis(5))
                                            .withMaxPartSize(1024 * 1024 * 1024).build())
                            .build();
    
                    resultSum
                            .map(new ShippingPriorityItemMap(pinningPolicy)).name(ShippingPriorityItemMap.class.getSimpleName()).uid(ShippingPriorityItemMap.class.getSimpleName())
                            .addSink(sink).name(OPERATOR_SINK).uid(OPERATOR_SINK);
                } else {
                    System.out.println("discarding output");
                }
    
                System.out.println("Stream job: " + TPCHQuery03.class.getSimpleName());
                System.out.println("Execution plan >>>\n" + env.getExecutionPlan());
                env.execute(TPCHQuery03.class.getSimpleName());
            } catch (IOException e) {
                e.printStackTrace();
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    
        public static void main(String[] args) throws Exception {
            new TPCHQuery03();
        }
    }
    

    UDF 在这里:OrderSource、OrderKeyedByCustomerProcessFunction、ShippingPriorityKeyedProcessFunction 和 SumShippingPriorityItem。我正在使用com.google.common.collect.ImmutableList,因为状态不会更新。此外,我只保留状态中必要的列,例如 ImmutableList&lt;Tuple2&lt;Long, Double&gt;&gt; lineItemList。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-12-04
      相关资源
      最近更新 更多