【问题标题】:Spark Streaming: How to process using multiple inputs to job?Spark Streaming:如何使用多个输入来处理作业?
【发布时间】:2016-03-01 16:26:25
【问题描述】:

输入 1: KV 数据流。
输入 2: 一些静态数据分区(用于处理输入 1 中的流)
问题可以建模为下图:

与 HDFS/RDD 分区共存: 我们如何确保流式任务 Map1、Map2 和 Map3 在以下机器上运行HDFS/RDD 分区是否存在?

图像描述:假设 K 是流式键(不是元组)。 First Map 将其转换为元组(具有空值)并将其广播到 3 个 Mapper。每个映射器都运行在包含不同分区的 RDD(或 HDFS 文件,这是第二个输入和静态数据)的不同节点上。每个 Mapper 使用 RDD 分区来计算键的值。最后我们要为键聚合值(使用reduceByKey _+_)。

【问题讨论】:

    标签: hadoop apache-spark stream spark-streaming flink-streaming


    【解决方案1】:

    如果我理解正确:

    1. K 是您从***DStream 通过streaming 工作获得的RDD。我不知道您传入数据的来源。这些数据基本上是一个array/seq/list 的 Keys。
    2. 您提到的静态数据是PairedRDD,格式为<K, Object>。从Object 中,您想为incoming RDD 中的键提取Val_n。
    3. 您的目标是在此join(或查找)过程中避免/最小化shuffle。

    为此,最好的策略是使用Join 操作与incoming RDD 和Static RDD 与两个RDD partitioned 使用相同的Partitioner。如果其中一个数据 RDD 比另一个小得多,您可以探索broadcasting 较小的那个。我最近在我的项目中尝试过这个,并在帖子中分享了经验:Random Partitioner behavior on the joined RDD

    编辑:由于您要处理您的密钥,K(假设 K=Set{K1, K2...Kn}),使用StaticRDD,在分区的位置,我建议采用如下方法。我没有检查语法和正确性,但你会明白的。

    val kRddBroadcastVar = .... // broadcasted variable 
    val keyValRDD = staticRDD.mapPartitions {   
           iter => transformKRddToTuple2Events(iter, kRddBroadcastVar )
         }
    
    def transformKRddToTuple2Events( iter: Iterator[Object], kRddBroadcastVar: List[KeyObjectType] ) : Iterator[(keyObjectType, valueObjectType )] {    
         val staticList = iter.toList
         val toReturn   = kRddBroadcastVar.map ( k => getKeyValue(k, staticList) )    
         toReturn.iterator
    }
    
    val outRdd = keyValRDD.reduceByKey( _ + _ )
    

    如果这有意义,请将此答案标记为已接受。

    【讨论】:

    • 嗨莫希特!谢谢回答。我试图解释编辑中的图像。实际上输入键 K 是要广播给所有映射器的,所以共同分区没有帮助。其次,我们的静态 RDD 可能不包含元组。但是为了评估输入键的值,我们想要遍历分区(内存中)。请提出建议。
    • 嗨莫希特!感谢您编辑您的答案。给了我很多见解。我试图广播一个可变的 scala 集合。然后在静态 RDD 上使用 mapPartition 来操作广播变量的本地副本。它刚刚奏效!所有节点都有不同的广播变量值。现在我将尝试共同分区输入流并使用广播变量而不是数据分区。很快就会更新你。谢谢!
    • 嗨莫希特!我想做的事情,无法按我计划的方式完成(最后评论)。使用相同的分区器,spark 不能确保不同 RDD 的相同分区的共同定位。加入可能会这样做,但一般不会这样做!我的静态数据不是 key-val。流数据中的每个键都需要使用整个数据的 2 次完整扫描来处理。因此,当静态数据被划分为多个分区(希望位于不同的节点上)时,我希望将每个键(来自流输入)广播到具有静态数据分区的所有节点并进行本地扫描。有什么建议吗?
    • 两个 Rdd 的相同分区器不能确保协同定位。即使对于Join,spark也不能保证。无论如何,由于您不打算使用连接(并且您的数据不是 对),因此无需使用任何分区器。回到你的问题,如果你已经广播了你的关键数据,它将作为一个简单的 java 变量在所有节点上可用。换句话说,您在所有节点上都有此变量的副本。现在,我建议 mapPartition 方法采用 transformKRddToTuple2Events 来利用静态和关键数据在每个分区产生所需的结果。
    • 如果您需要进一步说明,请告诉我。
    【解决方案2】:

    您的静态 RDD 是否足够小,可以缓存。在这种情况下,Spark 将尝试在这些节点上运行流式传输任务。但它不能保证。

    另外,如果参考数据很小,为什么不广播该数据集。

    我们一直在尝试解决与我们的数据存储 SnappyData (http://www.snappydata.io/) 中的首选位置有关的类似问题,其中数据位置是一等公民。

    【讨论】:

    • 感谢 Rishitesh!实际上,每个键都需要对 Mapper1、Mapper2 和 Mapper3 的数据分区进行多次扫描。如果我们有完整的数据作为广播变量,多次扫描可能会很耗时。在较小的数据(分区)上,多次扫描的计算密集度较低。实际上对于我的用例(研究工作)Apache Samza 更适合,但我也需要与其他框架进行比较。
    猜你喜欢
    • 2017-03-29
    • 2017-08-24
    • 1970-01-01
    • 2018-07-15
    • 2018-10-26
    • 2015-06-19
    • 2017-02-26
    • 2016-08-31
    • 1970-01-01
    相关资源
    最近更新 更多