代码之家  ›  专栏  ›  技术社区  ›  Mohamad Shaker

Spark-如何按键进行条件约简?

  •  0
  • Mohamad Shaker  · 技术社区  · 8 年前

    我有一个包含两列(键、值)的数据帧,如下所示:

    +------------+--------------------+
    |         key|               value|
    +------------+--------------------+
    |[sid2, sid5]|             value1 |
    |      [sid2]|             value2 |
    |      [sid6]|             value3 |
    +------------+--------------------+
    

    +------------+--------------------+
    |         key|               value|
    +------------+--------------------+
    |[sid2, sid5]|   [value1, value2] |
    |      [sid6]|             value3 |
    +------------+--------------------+
    

    我试图使用case类作为键wapper,并重写equals和hashCode函数,但没有成功( SPARK-2620 ).

    知道怎么做吗? 提前谢谢。

    更新-数据帧架构:

    root
     |-- id1: array (nullable = true)
     |    |-- element: string (containsNull = true)
     |-- events1: array (nullable = true)
     |    |-- element: struct (containsNull = true)
     |    |    |-- sid: string (nullable = true)
     |    |    |-- uid: string (nullable = true)
     |    |    |-- action: string (nullable = true)
     |    |    |-- touchPoint: string (nullable = true)
     |    |    |-- result: string (nullable = true)
     |    |    |-- timestamp: long (nullable = false)
     |    |    |-- url: string (nullable = true)
     |    |    |-- onlineId: long (nullable = false)
     |    |    |-- channel: string (nullable = true)
     |    |    |-- category: string (nullable = true)
     |    |    |-- clientId: long (nullable = false)
     |    |    |-- newUser: boolean (nullable = false)
     |    |    |-- userAgent: string (nullable = true)
     |    |    |-- group: string (nullable = true)
     |    |    |-- pageType: string (nullable = true)
     |    |    |-- clientIP: string (nullable = true)
    
    3 回复  |  直到 8 年前
        1
  •  1
  •   user9003280    8 年前

    这不能用 reduceByKey 因为问题定义不适合 byKey 转换。核心要求是密钥具有定义良好的标识,但这里的情况并非如此。

    考虑我们有密钥的数据集 [sid2, sid4, sid5] [sid2, sid3, sid5] . 在这种情况下,无法将对象唯一地分配给分区。重写哈希代码对您毫无帮助。

    更糟糕的是,一般情况下的问题是分布式的。考虑一组集合,例如对于每个集合,至少有一个其他集合具有非空交点。在这种情况下,所有值都应该合并到一个“集群”中。

    总的来说,如果没有相当严格的限制,这对Spark来说不是一个好问题,并且无法用basic解决 byKey公司 根本不需要转换。

    效率低下的解决方案(可能部分解决您的问题)是使用笛卡尔积:

    rdd.cartesian(rdd)
      .filter { case ((k1, _), (k2, _)) => intersects(v1, v2) }
      .map { case ((k, _), (_, v)) => (k, v) }
      .groupByKey
      .mapValues(_.flatten.toSet)
    

    然而,这是低效的,不能解决歧义。

        2
  •  1
  •   Jacek Laskowski    8 年前

    我认为,使用Spark SQL的数据集API是可行的(结果是直接翻译了基于RDD的@user9003280解决方案)。

    // the dataset
    val kvs = Seq(
      (Seq("sid2", "sid5"), "value1"),
      (Seq("sid2"), "value2"),
      (Seq("sid6"), "value3")).toDF("key", "value")
    scala> kvs.show
    +------------+------+
    |         key| value|
    +------------+------+
    |[sid2, sid5]|value1|
    |      [sid2]|value2|
    |      [sid6]|value3|
    +------------+------+
    
    val intersect = udf { (ss: Seq[String], ts: Seq[String]) => ss intersect ts }
    val solution = kvs.as("left")
      .join(kvs.as("right"))
      .where(size(intersect($"left.key", $"right.key")) > 0)
      .select($"left.key", $"right.value")
      .groupBy("key")
      .agg(collect_set("value") as "values")
      .dropDuplicates("values")
    scala> solution.show
    +------------+----------------+
    |         key|          values|
    +------------+----------------+
    |      [sid6]|        [value3]|
    |[sid2, sid5]|[value2, value1]|
    +------------+----------------+
    
        3
  •  0
  •   Mohamad Shaker    8 年前

    我在100000行数据帧上尝试了笛卡尔乘积解决方案,它花费了很多时间来处理,所以我决定使用graph GraphFrame ,可以直接在线性时间内计算图的连通分量(根据图的顶点和边的数量)。

    • 创建顶点和边数据帧。
    • 构建图表。

    最终结果如下:

    +------------+------+----------
    |         key| value|component
    +------------+------+----------
    |      [sid5]|value1|component1
    |      [sid2]|value2|component1
    |      [sid6]|value3|component2
    +------------+------+-----------
    

    然后是groupBy(“组件”)

    就是这样:)