【问题标题】:How to get Flink taskmanager number before job submitted?如何在提交作业之前获取 Flink taskmanager 编号?
【发布时间】:2021-06-17 14:37:23
【问题描述】:

我有一个 Flink Datastream 作业由

启动
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setParallelism(taskmanagernumber * x) // set env parallelism this line

env.addSource...map...addSink...
env.execute()

我想控制与 taskmanager 编号相关的 env 并行度,如上面的代码。

有办法吗?或者任何解决方法来设置与任务管理器编号相关的并行度?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    您可以使用reactive scheduler,它将自动调整并行度以适应集群提供的任何内容。

    您不必在作业本身中设置并行度。您可以在启动作业时在命令行中设置它:

    flink run -p <parallelism> <jar-file> <arguments>
    

    如果您不知道集群有多少可用插槽,您可以从 REST API 获取信息。 /overview 返回如下内容:

    {
      taskmanagers: 2,
      slots-total: 2,
      slots-available: 2,
      jobs-running: 0,
      jobs-finished: 0,
      jobs-cancelled: 0,
      jobs-failed: 0,
      flink-version: "1.13.1",
      flink-commit: "a7f3192"
    }
    

    slots-available 是您正在寻找的。所以你可以做类似的事情

    flink run -p `curl -s http://localhost:8081/overview | jq '.["slots-available"]'` ...
    

    【讨论】:

    • 如果我想在代码中设置呢?就像我有代码加载配置并将配置设置为 Flink,我认为命令行不是处理这种情况的好方法。
    • 总的来说,我认为最好将操作配置保留在代码之外。但您可以将值作为参数传递给作业,或从作业内部调用 REST API。你也可以在 flink-conf.yaml 中设置并行度。
    猜你喜欢
    • 2019-02-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-18
    • 1970-01-01
    • 1970-01-01
    • 2019-09-17
    相关资源
    最近更新 更多