• Flink 流处理API


    目录

    Environment

     Source

    从集合读取数据

    从文件读取数据

    以Kafka消息队列的数据作为来源


    流处理流程

    Environment

    getExecutionEnvironment

    创建一个执行环境,表示当前执行程序的上下文。如果程序是独立调用的,则此方法返回本地执行环境;如果从命令行客户端调用程序以提交到集群,则此方法返回此集群的执行环境,也就是说,getExecutionEnvironment会根据查询运行的方法决定返回什么样的运行环境,是常用的一种创建执行环境的方式。

    1. //创建一个执行环境
    2. val env=StreamExecutionEnvironment.getExecutionEnvironment

    如果没有设置并行度,会以 flink-conf.yaml 中的配置为准,默认是 1

     Source

    从集合读取数据

    测试代码

    1. package com.atguigu.apitest
    2. import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
    3. import org.apache.flink.streaming.api.scala._
    4. import java.util.Properties
    5. //
    6. case class SensorReading( id: String, timestamp: Long, temperature: Double)
    7. object SourceTest {
    8. def main(args: Array[String]): Unit = {
    9. //创建一个执行环境
    10. val env=StreamExecutionEnvironment.getExecutionEnvironment
    11. //从集合中获取数据
    12. val dataList = List(
    13. SensorReading("sensor_1", 1547718199, 35.8),
    14. SensorReading("sensor_6", 1547718201, 15.4),
    15. SensorReading("sensor_7", 1547718202, 6.7),
    16. SensorReading("sensor_10", 1547718205, 38.1)
    17. )
    18. val stream1=env.fromCollection(dataList)
    19. //打印输出
    20. stream1.print()
    21. //执行
    22. env.execute("source test")
    23. }
    24. }

    测试结果 

    从文件读取数据

    测试代码

    1. package com.atguigu.wc
    2. import org.apache.flink.api.scala.ExecutionEnvironment
    3. import org.apache.flink.api.scala._
    4. object WordCount {
    5. def main(args: Array[String]): Unit = {
    6. //创建一个批处理的执行环境
    7. val env:ExecutionEnvironment = ExecutionEnvironment.getExecutionEnvironment
    8. //从文件中读取数据
    9. val inputPath = "D:\\HYF\\FlinkTutorial\\src\\main\\resources\\sensor.txt"
    10. val stream2 = env.readTextFile(inputPath)
    11. //打印输出
    12. stream2.print()
    13. //执行
    14. env.execute("source test")
    15. }
    16. }

    测试结果

    Kafka消息队列的数据作为来源

    需要pom.xml文件中引入 kafka 连接器的依赖:

    1. <dependency>
    2. <groupId>org.apache.flink</groupId>
    3. <artifactId>flink-connector-kafka-0.11_2.12</artifactId>
    4. <version>1.10.1</version>
    5. </dependency>

    测试环境

    1、先打开zookeeper

    2、再打开Kafka

    3、jpsall查看是否启动成功

    4、先完全关闭Kafka才能关闭zookeeper

    测试代码

    1. package com.atguigu.apitest
    2. import org.apache.flink.api.common.serialization.SimpleStringSchema
    3. import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
    4. import org.apache.flink.streaming.api.scala._
    5. import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer011
    6. import java.util.Properties
    7. case class SensorReading( id: String, timestamp: Long, temperature: Double)
    8. object SourceTest {
    9. def main(args: Array[String]): Unit = {
    10. //创建一个执行环境
    11. val env=StreamExecutionEnvironment.getExecutionEnvironment
    12. //从集合中获取数据
    13. val dataList = List(
    14. SensorReading("sensor_1", 1547718199, 35.8),
    15. SensorReading("sensor_6", 1547718201, 15.4),
    16. SensorReading("sensor_7", 1547718202, 6.7),
    17. SensorReading("sensor_10", 1547718205, 38.1)
    18. )
    19. val stream1=env.fromCollection(dataList)
    20. //打印输出
    21. //stream1.print()
    22. //从文件中读取数据
    23. val inputPath = "D:\\HYF\\FlinkTutorial\\src\\main\\resources\\sensor.txt"
    24. val stream2 = env.readTextFile(inputPath)
    25. //打印输出
    26. //stream2.print()
    27. //从kafka读取数据
    28. val properties = new Properties()
    29. properties.setProperty("bootstrap.servers","hadoop102:9092")
    30. properties.setProperty("group.id","consumer-group")
    31. val stream3 = env.addSource( new FlinkKafkaConsumer011[String]("sensor",new
    32. SimpleStringSchema(),properties) )
    33. stream3.print()
    34. //执行
    35. env.execute("source test")
    36. }
    37. }

    测试结果

  • 相关阅读:
    Java题目集-Chapter 10 Object-Oriented Thinking
    使用 Data Assistant 快速创建测试数据集
    Python数据分析与机器学习4-Matplotlib
    磁盘的工作方式
    洛谷P2451 遗传代码
    中国SSD产业突围有多难?除了技术“瓶颈”还有哪里挑战?
    【C++】静态库lib和动态库dll的优缺点、使用方法
    使用MASA Blazor开发一个标准的查询表格页
    【论文笔记】(FGSM)Explaining and Harnessing Adversarial Examples
    matcher中find,matches,lookingAt匹配字符串的不同之处说明
  • 原文地址:https://blog.csdn.net/qq_70085330/article/details/126614462