Flink 在 Idea上提交任务到远程服务器

Flink自身提供了远程提交任务的环境,源码如下:

请查看StreamExecutionEnvironment 类中 createRemoteEnvironment 方法

def createRemoteEnvironment(
    host: String,
    port: Int,
    parallelism: Int,
    jarFiles: String*): StreamExecutionEnvironment = {

  val javaEnv = JavaEnv.createRemoteEnvironment(host, port, jarFiles: _*)
  javaEnv.setParallelism(parallelism)
  new StreamExecutionEnvironment(javaEnv)
}
远程提交示例代码如下:
package com.flink.remotesubmit

import org.apache.flink.streaming.api.scala._

object RemoteSubmitApp extends App {

  val host: String = "node02"
  val port: Int = 8081
  val jarFiles = "E:\\CDHProjectDemo\\flink-demo\\target\\flink-demo-0.0.1-SNAPSHOT.jar"

  val env = StreamExecutionEnvironment.createRemoteEnvironment(host, port, jarFiles)

  val socketHost: String = "node01"
  val socketPort: Int = 7777
  val socketDs: DataStream[String] = env.socketTextStream(socketHost, socketPort)

  socketDs.flatMap(_.split(" "))
    .map((_, 1))
    .keyBy(0)
    .sum(1)
    .print()

  env.execute("Remote Submit Job")
}

注意:

  1. 需要保持代码和jar一致性,意思就是修改代码之后需重新执行 mvn clean package
  2. 需在项目的 src/main/resource 目录中添加相关配置文件(core-site.xml、hdfs-site.xml、yarn-site.xml、mapred-site.xml等)
运行情况
image
image
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

友情链接更多精彩内容