1 在 idea 中添加依赖

1.1 在创建目录时添加

  1. 使用 Intellij IDEA 创建一个 Maven 新项目
  2. 勾选 Create from archetype,然后点击 Add Archetype 按钮
  3. GroupId 中输入 org.apache.flink, ArtifactId 中输入 flink-quickstart-scala,Version 中输入 1.10.0,然后点击 OK
  4. 点击向右箭头,出现下拉列表,选中 flink-quickstart-scala:1.10.0,点击 Next
  5. Name 中输入 FlinkTutorial, GroupId 中输入 com.atguigu, ArtifactId 中输入FlinkTutorial,点击 Next
  6. 最好使用 IDEA 默认的 Maven 工具Bundled(Maven 3),点击 Finish,等待一会儿,项目就创建好了

2 编写代码

如下编写 wordcount 代码:

package test1

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

object WordCountFromBatch {
  def main(args: Array[String]): Unit = {
    // 获取运行时环境,类似SparkContext
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    // 并行任务的数量设置为1
    // 全局并行度
    env.setParallelism(1)

    val stream = env
      .fromElements(
        "xiamao",
        "hello world",
        "xiaomao xiaogou",
        "xiaozhu hello"
      )
      .setParallelism(1)

    // 对数据流进行转换算子操作
    val textStream = stream
      // 使用空格来进行切割输入流中的字符串
      .flatMap(r => r.split("\\s"))
      .setParallelism(2)
      // 做map操作, w => (w, 1)
      .map(w => WordWithCount(w, 1))
      .setParallelism(2)
      // 使用word字段进行分组操作,也就是shuffle
      .keyBy(0)
      // 做聚合操作,类似与reduce
      .sum(1)
        .setParallelism(2)

    // 将数据流输出到标准输出,也就是打印
    // 设置并行度为1,print算子的并行度就是1,覆盖了全局并行度
    textStream.print().setParallelism(2)

    // 不要忘记执行!
    env.execute()
  }

  case class WordWithCount(word: String, count: Int)
}

2.1 创建运行时环境

val env = StreamExecutionEnvironment.getExecutionEnvironment

注意:

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

这个要导入该包所有的东西,因为有好多隐式转换要用.

2.2 添加 source生成流

env
      .fromElements(
        "xiamao",
        "hello world",
        "xiaomao xiaogou",
        "xiaozhu hello"
      )

这里用的是静态元素来生成流.

2.2.1 fromElements方法专门用几台数据生成流.

.fromElements((1, 1L))

2.2.2 也可以用socket来生成流:

val stream = env.socketTextStream("localhost", 9999, '\n')

还可以使用addSource接口来添加其他的流.

2.3 计算

根据业务逻辑计算.

2.4 添加sink

就是计算完的数据发往哪里.
我们这里直接输出,相当于sink是控制台.

可以通过addSink来设置Sink的类型.

2.5 执行程序

env.execute()

2.6 编译执行或者打成jar包

3 jar在flink上运行

3.1 启动 Flink 集群

$ cd flink-1.10.0
$ ./bin/start-cluster.sh

可以打开 Flink WebUI 查看集群状态:

http://localhost:8081

3.2 提交打包好的 JAR 包

$ cd flink-1.10.0
$ ./bin/flink run 打包好的 JAR 包的绝对路径

也可以在webUI上提交

3.3 停止 Flink 集群

$ ./bin/stop-cluster.sh

3.4 查看标准输出日志的位置,在 log 文件夹中

$ cd flink-1.10.0/log
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐