日韩av黄I国产麻豆传媒I国产91av视频在线观看I日韩一区二区三区在线看I美女国产在线I麻豆视频国产在线观看I成人黄色短片

歡迎訪問 生活随笔!

生活随笔

當(dāng)前位置: 首頁 >

Scala模拟Spark分布式计算流程示例代码

發(fā)布時(shí)間:2025/1/21 33 豆豆
生活随笔 收集整理的這篇文章主要介紹了 Scala模拟Spark分布式计算流程示例代码 小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,幫大家做個(gè)參考.

場景

兩個(gè)Executor,分別接收來自Driver分發(fā)的任務(wù)(數(shù)據(jù)和計(jì)算邏輯)

代碼

Executor1

package com.zxl.bigdata.spark.core.testimport java.io.{InputStream, ObjectInputStream} import java.net.{ServerSocket, Socket}object Executor {def main(args: Array[String]): Unit = {// 啟動(dòng)服務(wù)器,接收數(shù)據(jù)val server = new ServerSocket(9999)println("服務(wù)器啟動(dòng),等待接收數(shù)據(jù)")// 等待客戶端的連接val client: Socket = server.accept()val in: InputStream = client.getInputStreamval objIn = new ObjectInputStream(in)val task: SubTask = objIn.readObject().asInstanceOf[SubTask]val ints: List[Int] = task.compute()println("計(jì)算節(jié)點(diǎn)[9999]計(jì)算的結(jié)果為:" + ints)objIn.close()client.close()server.close()} }

Executor2

package com.zxl.bigdata.spark.core.testimport java.io.{InputStream, ObjectInputStream} import java.net.{ServerSocket, Socket}object Executor2 {def main(args: Array[String]): Unit = {// 啟動(dòng)服務(wù)器,接收數(shù)據(jù)val server = new ServerSocket(8888)println("服務(wù)器啟動(dòng),等待接收數(shù)據(jù)")// 等待客戶端的連接val client: Socket = server.accept()val in: InputStream = client.getInputStreamval objIn = new ObjectInputStream(in)val task: SubTask = objIn.readObject().asInstanceOf[SubTask]val ints: List[Int] = task.compute()println("計(jì)算節(jié)點(diǎn)[8888]計(jì)算的結(jié)果為:" + ints)objIn.close()client.close()server.close()} }

Task

package com.zxl.bigdata.spark.core.testclass Task extends Serializable {val datas = List(1,2,3,4)//val logic = ( num:Int )=>{ num * 2 }val logic : (Int)=>Int = _ * 2}

SubTask

package com.zxl.bigdata.spark.core.testclass SubTask extends Serializable {var datas : List[Int] = _var logic : (Int)=>Int = _// 計(jì)算def compute() = {datas.map(logic)} }

Driver

package com.zxl.bigdata.spark.core.testimport java.io.{ObjectOutputStream, OutputStream} import java.net.Socketobject Driver {def main(args: Array[String]): Unit = {// 連接服務(wù)器val client1 = new Socket("localhost", 9999)val client2 = new Socket("localhost", 8888)val task = new Task()val out1: OutputStream = client1.getOutputStreamval objOut1 = new ObjectOutputStream(out1)val subTask = new SubTask()subTask.logic = task.logicsubTask.datas = task.datas.take(2)objOut1.writeObject(subTask)objOut1.flush()objOut1.close()client1.close()val out2: OutputStream = client2.getOutputStreamval objOut2 = new ObjectOutputStream(out2)val subTask1 = new SubTask()subTask1.logic = task.logicsubTask1.datas = task.datas.takeRight(2)objOut2.writeObject(subTask1)objOut2.flush()objOut2.close()client2.close()println("客戶端數(shù)據(jù)發(fā)送完畢")} }

程序運(yùn)行日志


總結(jié)

以上是生活随笔為你收集整理的Scala模拟Spark分布式计算流程示例代码的全部內(nèi)容,希望文章能夠幫你解決所遇到的問題。

如果覺得生活随笔網(wǎng)站內(nèi)容還不錯(cuò),歡迎將生活随笔推薦給好友。