IT数码 购物 网址 头条 软件 日历 阅读 图书馆
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
图片批量下载器
↓批量下载图片,美女图库↓
图片自动播放器
↓图片自动播放器↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁
 
   -> 大数据 -> 07-flink-1.10.1- 用户自定义 flink source api -> 正文阅读

[大数据]07-flink-1.10.1- 用户自定义 flink source api

1 定义自己想要的流返回值数据类型

package com.study.liucf.unbounded.source

/**
 * @Author liucf
 * @Date 2021/9/8
 */
case class LiucfSensorReding(id:String,timestamp:Long,temperature:Double)

2 定义自己的读取数据源的类

package com.study.liucf.unbounded.source

import java.util.Random

import org.apache.flink.streaming.api.functions.source.SourceFunction


/**
 * @Author liucf
 * @Date 2021/9/8
 *
 *      我读取数据源类
 */
class LiucfSourceFunction() extends SourceFunction[LiucfSensorReding]{
  var running = true
  //模拟读取数据
  override def run(ctx: SourceFunction.SourceContext[LiucfSensorReding]): Unit = {
    //定义随机变量
    val random = new Random()
    //定义一组10个传感器初始温度
    var currentTemp = 1.to(10).map(r=>("sensor_"+r,random.nextDouble()*100))
    while(running){
      currentTemp = currentTemp.map(r=>(r._1,r._2+random.nextGaussian()))
      val currentTimestamp = System.currentTimeMillis()
      //逐条发出去
      currentTemp.foreach(r=>ctx.collect(LiucfSensorReding(r._1,currentTimestamp,r._2)))
    }
    Thread.sleep(1000)
  }

  //如果调用cancel方法,则running标准位会被设置成false,停止读取数据
  override def cancel(): Unit = {
     running = false
  }
}

3 flink使用自定义的数据源处理数据

package com.study.liucf.unbounded.source

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

/**
 * @Author liucf
 * @Date 2021/9/8
 *
 *      自定义数据源读取数据
 */
object MySource {
  def main(args: Array[String]): Unit = {
    //创建flink执行环境
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    //添加自定义数据源
    val ds: DataStream[LiucfSensorReding] = env.addSource(new LiucfSourceFunction())
    //输出到标准控制台
    ds.print()
    //启动flink运行
    env.execute("liucf defined source api")

  }

}

输出结果

?

  大数据 最新文章
实现Kafka至少消费一次
亚马逊云科技:还在苦于ETL?Zero ETL的时代
初探MapReduce
【SpringBoot框架篇】32.基于注解+redis实现
Elasticsearch:如何减少 Elasticsearch 集
Go redis操作
Redis面试题
专题五 Redis高并发场景
基于GBase8s和Calcite的多数据源查询
Redis——底层数据结构原理
上一篇文章      下一篇文章      查看所有文章
加:2021-09-09 11:50:45  更:2021-09-09 11:51:58 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2025年1日历 -2025/1/18 14:33:26-

图片自动播放器
↓图片自动播放器↓
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
图片批量下载器
↓批量下载图片,美女图库↓
  网站联系: qq:121756557 email:121756557@qq.com  IT数码