RDD的join和Dstream的join有什么区别?

2023-10-09 03:32
文章标签 区别 join rdd dstream

本文主要是介绍RDD的join和Dstream的join有什么区别?,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

有人在知识星球里问:

浪院长,RDD的join和Dstream的join有什么区别?

浪尖的回答:

DStream的join底层就是rdd的join。

下面,我们就带着疑问去验证以下,我们的想法。

2. DStream -> PairDStreamFunctions

Dstream这个类实际上支持的只是Spark Streaming的基础操作算子,比如: mapfilter 和window.PairDStreamFunctions 这个支持key-valued类型的流数据

,支持的操作算子,如,groupByKeyAndWindow,join。这些操作,在有key-value类型的流上是自动识别的。

对于dstream -> PairDStreamFunctions自动转换的过程大家肯定想到的是scala的隐式转换。具体代码在Dstream的object内部。

implicit def toPairDStreamFunctions[K, V](stream: DStream[(K, V)])
     (implicit kt: ClassTag[K], vt: ClassTag[V], ord: Ordering[K] = null):
   PairDStreamFunctions[K, V] = {
   new PairDStreamFunctions[K, V](stream)
 }

假如,你对scala的隐式转换比较懵逼,请阅读下面文章。

Scala语法基础之隐式转换

3. PairDStreamFunctions的join

PairDStreamFunctions的join API总共有三种

/**
  * Return a new DStream by applying 'join' between RDDs of `this` DStream and `other` DStream.
  * Hash partitioning is used to generate the RDDs with Spark's default number of partitions.
   *
   *  通过join this和other Dstream的rdd构建出一个新的DStream.
   *  Hash分区器,用来使用默认的分区数来产生RDDs。
  */
 def join[W: ClassTag](other: DStream[(K, W)]): DStream[(K, (V, W))] = ssc.withScope {
   join[W](other, defaultPartitioner())
 }

 /**
  * Return a new DStream by applying 'join' between RDDs of `this` DStream and `other` DStream.
  * Hash partitioning is used to generate the RDDs with `numPartitions` partitions.
   *
   *  通过join this和other Dstream的rdd构建出一个新的DStream.
   *  Hash分区器,用来使用numPartitions分区数来产生RDDs。
  */
 def join[W: ClassTag](
     other: DStream[(K, W)],
     numPartitions: Int): DStream[(K, (V, W))] = ssc.withScope {
   join[W](other, defaultPartitioner(numPartitions))
 }

 /**
  * Return a new DStream by applying 'join' between RDDs of `this` DStream and `other` DStream.
  * The supplied org.apache.spark.Partitioner is used to control the partitioning of each RDD.
   * 通过join this和other Dstream的rdd构建出一个新的DStream.
   * 使用org.apache.spark.Partitioner来控制每个RDD的分区。
  */
 def join[W: ClassTag](
     other: DStream[(K, W)],
     partitioner: Partitioner
   ): DStream[(K, (V, W))] = ssc.withScope {
   self.transformWith(
     other,
     (rdd1: RDD[(K, V)], rdd2: RDD[(K, W)]) => rdd1.join(rdd2, partitioner)
   )
 }

上面所示代码中,第三个PairDStreamFunctions的join api 体现了join的操作,也即是函数:

(rdd1: RDD[(K, V)], rdd2: RDD[(K, W)]) => rdd1.join(rdd2, partitioner)

上面是两个RDD的join过程,并且指定了分区器。后面我们主要是关注该函数封装及调用。

其实,看过浪尖的Spark Streaming的视频的朋友或者度过浪尖关于Spark Streaming相关源码讲解的朋友应该有所了解的是。 这个生成RDD的函数应该是在 DStream的compute方法中在生成RDD的时候调用。假设你不了解也不要紧。 我们跟着代码轨迹前进,验证我们的想法。

DStream.transformWith

/**
  * Return a new DStream in which each RDD is generated by applying a function
  * on each RDD of 'this' DStream and 'other' DStream.
  */

 def transformWith[U: ClassTag, V: ClassTag](
     other: DStream[U], transformFunc: (RDD[T], RDD[U]) => RDD[V]
   ): DStream[V] = ssc.withScope {
   // because the DStream is reachable from the outer object here, and because
   // DStreams can't be serialized with closures, we can't proactively check
   // it for serializability and so we pass the optional false to SparkContext.clean
   val cleanedF = ssc.sparkContext.clean(transformFunc, false)
   transformWith(other, (rdd1: RDD[T], rdd2: RDD[U], time: Time) => cleanedF(rdd1, rdd2))
 }
 
   /**
    * Return a new DStream in which each RDD is generated by applying a function
    * on each RDD of 'this' DStream and 'other' DStream.
    */

   def transformWith[U: ClassTag, V: ClassTag](
       other: DStream[U], transformFunc: (RDD[T], RDD[U], Time) => RDD[V]
     ): DStream[V] = ssc.withScope {
     // because the DStream is reachable from the outer object here, and because
     // DStreams can't be serialized with closures, we can't proactively check
     // it for serializability and so we pass the optional false to SparkContext.clean
     val cleanedF = ssc.sparkContext.clean(transformFunc, false)
     val realTransformFunc = (rdds: Seq[RDD[_]], time: Time) => {
       assert(rdds.length == 2)
       val rdd1 = rdds(0).asInstanceOf[RDD[T]]
       val rdd2 = rdds(1).asInstanceOf[RDD[U]]
       cleanedF(rdd1, rdd2, time)
     }
     new TransformedDStream[V](Seq(this, other), realTransformFunc)
   }

经过上面两个 TransformWith操作,最终生成了一个TransformedDStream。需要关注的是new TransformedDStream[V](Seq(this, other), realTransformFunc) 第一个参数是一个包含要进行join操作的两个流的Seq。

那么,TransformedDStream 的parents 就包含了两个流。我们可以看到其 compute 方法的第一行。

override def compute(validTime: Time): Option[RDD[U]] = {
//    针对每一个流,获取其当前时间的RDD。
   val parentRDDs = parents.map { parent => parent.getOrCompute(validTime).getOrElse(
     // Guard out against parent DStream that return None instead of Some(rdd) to avoid NPE
     throw new SparkException(s"Couldn't generate RDD from parent at time $validTime"))
   }

compute的第一行就是获取parent中每个流,当前有效时间的RDD。然后调用,前面步骤封装的函数进行join。

val transformedRDD = transformFunc(parentRDDs, validTime)

以上就是join的全部过程。也是,验证浪尖所说的,DStream的join底层就是RDD的join。

640?wx_fmt=jpeg

这篇关于RDD的join和Dstream的join有什么区别?的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/170049

相关文章

go 指针接收者和值接收者的区别小结

《go指针接收者和值接收者的区别小结》在Go语言中,值接收者和指针接收者是方法定义中的两种接收者类型,本文主要介绍了go指针接收者和值接收者的区别小结,文中通过示例代码介绍的非常详细,需要的朋友们下... 目录go 指针接收者和值接收者的区别易错点辨析go 指针接收者和值接收者的区别指针接收者和值接收者的

售价599元起! 华为路由器X1/Pro发布 配置与区别一览

《售价599元起!华为路由器X1/Pro发布配置与区别一览》华为路由器X1/Pro发布,有朋友留言问华为路由X1和X1Pro怎么选择,关于这个问题,本期图文将对这二款路由器做了期参数对比,大家看... 华为路由 X1 系列已经正式发布并开启预售,将在 4 月 25 日 10:08 正式开售,两款产品分别为华

MySQL高级查询之JOIN、子查询、窗口函数实际案例

《MySQL高级查询之JOIN、子查询、窗口函数实际案例》:本文主要介绍MySQL高级查询之JOIN、子查询、窗口函数实际案例的相关资料,JOIN用于多表关联查询,子查询用于数据筛选和过滤,窗口函... 目录前言1. JOIN(连接查询)1.1 内连接(INNER JOIN)1.2 左连接(LEFT JOI

kotlin中const 和val的区别及使用场景分析

《kotlin中const和val的区别及使用场景分析》在Kotlin中,const和val都是用来声明常量的,但它们的使用场景和功能有所不同,下面给大家介绍kotlin中const和val的区别,... 目录kotlin中const 和val的区别1. val:2. const:二 代码示例1 Java

CSS Padding 和 Margin 区别全解析

《CSSPadding和Margin区别全解析》CSS中的padding和margin是两个非常基础且重要的属性,它们用于控制元素周围的空白区域,本文将详细介绍padding和... 目录css Padding 和 Margin 全解析1. Padding: 内边距2. Margin: 外边距3. Padd

Springboot @Autowired和@Resource的区别解析

《Springboot@Autowired和@Resource的区别解析》@Resource是JDK提供的注解,只是Spring在实现上提供了这个注解的功能支持,本文给大家介绍Springboot@... 目录【一】定义【1】@Autowired【2】@Resource【二】区别【1】包含的属性不同【2】@

Java中的String.valueOf()和toString()方法区别小结

《Java中的String.valueOf()和toString()方法区别小结》字符串操作是开发者日常编程任务中不可或缺的一部分,转换为字符串是一种常见需求,其中最常见的就是String.value... 目录String.valueOf()方法方法定义方法实现使用示例使用场景toString()方法方法

分辨率三兄弟LPI、DPI 和 PPI有什么区别? 搞清分辨率的那些事儿

《分辨率三兄弟LPI、DPI和PPI有什么区别?搞清分辨率的那些事儿》分辨率这个东西,真的是让人又爱又恨,为了搞清楚它,我可是翻阅了不少资料,最后发现“小7的背包”的解释最让我茅塞顿开,于是,我... 在谈到分辨率时,我们经常会遇到三个相似的缩写:PPI、DPI 和 LPI。虽然它们看起来差不多,但实际应用

GORM中Model和Table的区别及使用

《GORM中Model和Table的区别及使用》Model和Table是两种与数据库表交互的核心方法,但它们的用途和行为存在著差异,本文主要介绍了GORM中Model和Table的区别及使用,具有一... 目录1. Model 的作用与特点1.1 核心用途1.2 行为特点1.3 示例China编程代码2. Tab

Nginx指令add_header和proxy_set_header的区别及说明

《Nginx指令add_header和proxy_set_header的区别及说明》:本文主要介绍Nginx指令add_header和proxy_set_header的区别及说明,具有很好的参考价... 目录Nginx指令add_header和proxy_set_header区别如何理解反向代理?proxy