【问题标题】:Apache Flink - custom java options are not recognized inside jobApache Flink - 在作业中无法识别自定义 Java 选项
【发布时间】:2017-02-20 12:27:59
【问题描述】:

我在 flink-conf.yaml 中添加了以下行:

env.java.opts: "-Ddy.props.path=/PATH/TO/PROPS/FILE"

当启动 jobmanager (jobmanager.sh start cluster) 我在日志中看到 jvm 选项确实被识别

2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -  JVM Options:
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -Xms256m
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -Xmx256m
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -XX:MaxPermSize=256m
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -Ddy.props.path=/srv/dy/stream-aggregators/aggregators.conf
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -Dlog.file=/srv/flink-1.2.0/log/flink-flink-jobmanager-0-flinkvm-master.log
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -Dlog4j.configuration=file:/srv/flink-1.2.0/conf/log4j.properties
2017-02-20 12:19:23,536 INFO  org.apache.flink.runtime.jobmanager.JobManager                -     -Dlogback.configurationFile=file:/srv/flink-1.2.0/conf/logback.xml

但是当我运行 flink 作业(flink run -d PROG.JAR)时,System.getProperty("dy.props.path") 返回 null(并且在打印系统属性时,我看到它确实不存在。 )

真正的问题是——如何设置在 flink-job 代码中可用的系统属性?

【问题讨论】:

    标签: java apache-flink flink-streaming


    【解决方案1】:

    这个问题和Flink[1]的运行时架构有很大关系。

    我了解到您在独立集群中运行您的工作。请记住,JobManagerTaskManagers 在不同的 jvm 实例中运行。您必须考虑每个代码块将在哪里执行。

    例如,mapfilter 等转换中的代码在 TaskManager 上执行。 您的入口类的main 方法中的代码在命令行工具flink 中执行,该工具当然没有设置系统属性,因为它会生成一个临时(-d)jvm 仅用于提交作业。

    如果您通过 WebUI 提交作业,则您的 main 方法中的代码将在 JobManager 上执行,因此将设置该属性。

    一般来说,我宁愿不鼓励通过系统属性传递程序参数,因为这是一种不好的做法。


    下面有一个简单的例子:

    我开始了:

    • JobManagerenv.java.opts:"-Ddy.props.path=jobmanager"
    • TaskManagerenv.java.opts:"-Ddy.props.path=taskmanager"

    我的工作代码如下:

    object Main {
      def main(args: Array[String]): Unit = {
        val env = StreamExecutionEnvironment.getExecutionEnvironment
        val stream = env.fromCollection(1 to 4)
    
        val prop = System.getProperty("dy.props.path")
        stream.map(_ => System.getProperty("dy.props.path") + "  mainArg: " + prop).print()
    
        env.execute("stream")
      }
    }
    

    当我通过flink工具提交代码时,输​​出如下:

    taskmanager  mainArg: null
    taskmanager  mainArg: null
    taskmanager  mainArg: null
    taskmanager  mainArg: null
    

    当它通过WebUI 提交时,我得到:

    taskmanager  mainArg: jobmanager
    taskmanager  mainArg: jobmanager
    taskmanager  mainArg: jobmanager
    taskmanager  mainArg: jobmanager
    

    【讨论】:

    • 那么如何通过带有 jvm 选项的命令行向 fork jvm 提交作业?我看到了 yarn 的选项,但是独立集群呢?
    • 在独立集群中,必须在提交作业之前预先生成任务管理器。您可以在启动任务管理器时设置 jvm 选项。
    猜你喜欢
    • 2022-10-12
    • 2013-07-01
    • 1970-01-01
    • 2015-08-05
    • 1970-01-01
    • 1970-01-01
    • 2019-08-27
    • 2016-01-30
    • 2012-08-04
    相关资源
    最近更新 更多