当前位置:   article > 正文

Flink第一更——Source源接入_filesource

filesource

1)从Collection接入数据

/**
  * 从集合中采集数据
  */
object FromCollection {
  def main(args: Array[String]): Unit = {
    // 1. 环境
    val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

    // 2. source  集合
    val stream: DataStream[SensorReading] = env.fromCollection(List(
      SensorReading("sensor_1", 1547718199, 35.80018327300259),
      SensorReading("sensor_6", 1547718201, 15.402984393403084),
      SensorReading("sensor_7", 1547718202, 6.720945201171228),
      SensorReading("sensor_10", 1547718205, 38.101067604893444)
    ))

    // 3. transformation

    // 4. sink
    stream.print("stream").setParallelism(1)

    // 5. execute
    env.execute("API Test")
  }
}

/**
  *
  * @param id          传感器ID
  * @param timeStamp   时间戳
  * @param temperature 传感器温度
  */
case class SensorReading(id: String, timeStamp: Long, temperature: Double)
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30
  • 31
  • 32
  • 33

2)从Text文本接入数据

object FromText {
  def main(args: Array[String]): Unit = {
    // 1. 环境
    val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

    // 2. source  Text
    val stream: DataStream[String] = env.readTextFile("data/words")

    // 3. transformation

    // 4. sink
    stream.print("stream").setParallelism(1)

    // 5. execute
    env.execute("API Test")
  }
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17

3)从Kafka接入数据

/**
  * 从集合中采集数据
  */
object FromCollection {
  def main(args: Array[String]): Unit = {
    // 1. 环境
    val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

    // 2. source  集合
    val stream: DataStream[SensorReading] = env.fromCollection(List(
      SensorReading("sensor_1", 1547718199, 35.80018327300259),
      SensorReading("sensor_6", 1547718201, 15.402984393403084),
      SensorReading("sensor_7", 1547718202, 6.720945201171228),
      SensorReading("sensor_10", 1547718205, 38.101067604893444)
    ))

    // 3. transformation

    // 4. sink
    stream.print("stream").setParallelism(1)

    // 5. execute
    env.execute("API Test")
  }
}

/**
  *
  * @param id          传感器ID
  * @param timeStamp   时间戳
  * @param temperature 传感器温度
  */
case class SensorReading(id: String, timeStamp: Long, temperature: Double)
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30
  • 31
  • 32
  • 33
声明:本文内容由网友自发贡献,不代表【wpsshop博客】立场,版权归原作者所有,本站不承担相应法律责任。如您发现有侵权的内容,请联系我们。转载请注明出处:https://www.wpsshop.cn/w/代码探险家/article/detail/886359
推荐阅读
相关标签
  

闽ICP备14008679号