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 小米 华为 单反 装机 图拉丁
 
   -> 嵌入式 -> 创建flink的MqttSource (scala版本) -> 正文阅读

[嵌入式]创建flink的MqttSource (scala版本)

一.编写 MqttClient类

class MqttClient {

  var url:String="tcp://192.168.174.206:1883"
  var flag:Boolean=true
  private val topicArr: Array[Topic] = List(new Topic("topic_test",QoS.AT_LEAST_ONCE),new Topic("topic_test",QoS.AT_MOST_ONCE),new Topic("topic_test",QoS.EXACTLY_ONCE)).toArray

  //开始连接
  def  start():FutureConnection={
     val mqtt = new MQTT()

    mqtt.setHost(url) //设置连接地址
    mqtt.setCleanSession(true)//是否清空会话信息
    mqtt.setReconnectAttemptsMax(10) //设置最大的重新连接次数
    mqtt.setReconnectDelay(10) //设置重新连接的时长
    mqtt.setKeepAlive(30) //设置心跳时间
    mqtt.setSendBufferSize(64) //设置缓冲区的大小
    mqtt.setClientId("clientid_0001")

    val connection: FutureConnection = mqtt.futureConnection()
    connection.connect()
    connection.subscribe(topicArr)

    return connection

  }

}

?2.编辑MqttSource

object MqttSource extends RichSourceFunction[String] {

   var flag:Boolean=true
  override def run(sourceContext: SourceFunction.SourceContext[String]): Unit = {
    val client = new MqttClient()
    val connection: FutureConnection = client.start()

    var msg:String=""

    while (true){
      val futureMessage: Future[Message] = connection.receive()
      val message: Message = futureMessage.await()
      msg = String.valueOf(message.getPayloadBuffer()).substring(6)
      sourceContext.collect(msg)
    }
  }
  override def cancel(): Unit = {
    flag=false
  }
}

?三,编辑flink主类?flinkMQtt

object flinkMQtt {
  def main(args: Array[String]): Unit = {
  //创建执行的环境
   val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment
  //设置并行度
   env.setParallelism(1)
   val mqttDS: DataStreamSource[String] = env.addSource(MqttSource)
   mqttDS.print()
    
//执行job
env.execute()
  }
}

?存在一个问题:如果发送汉字字符串接,接收端收到一段数字?还没解决!!!知道的留下解决方式......

  嵌入式 最新文章
基于高精度单片机开发红外测温仪方案
89C51单片机与DAC0832
基于51单片机宠物自动投料喂食器控制系统仿
《痞子衡嵌入式半月刊》 第 68 期
多思计组实验实验七 简单模型机实验
CSC7720
启明智显分享| ESP32学习笔记参考--PWM(脉冲
STM32初探
STM32 总结
【STM32】CubeMX例程四---定时器中断(附工
上一篇文章      下一篇文章      查看所有文章
加:2021-09-04 17:42:22  更:2021-09-04 17:42:52 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2024年11日历 -2024/11/26 0:47:22-

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