【问题标题】:Does Spark Dataframe have an equivalent option of Panda's merge indicator?Spark Dataframe 是否具有 Panda 合并指示器的等效选项?
【发布时间】:2016-12-07 20:28:17
【问题描述】:

python Pandas 库包含以下函数:

DataFrame.merge(right, how='inner', on=None, left_on=None, right_on=None, left_index=False,
                right_index=False, sort=False, suffixes=('_x', '_y'), copy=True,
                indicator=False)

指标字段结合 Panda 的 value_counts() 函数可用于快速确定连接的执行情况。

例子:

In [48]: df1 = pd.DataFrame({'col1': [0, 1], 'col_left':['a', 'b']})

In [49]: df2 = pd.DataFrame({'col1': [1, 2, 2],'col_right':[2, 2, 2]})

In [50]: pd.merge(df1, df2, on='col1', how='outer', indicator=True)
Out[50]: 
   col1 col_left  col_right      _merge
0     0        a        NaN   left_only
1     1        b        2.0        both
2     2      NaN        2.0  right_only
3     2      NaN        2.0  right_only

在 Spark Dataframe 中检查连接性能的最佳方法是什么?

在 1 个答案中提供了一个自定义函数:它还没有给出正确的结果,但如果它可以的话那就太好了:

ASchema = StructType([StructField('id', IntegerType(),nullable=False),
                 StructField('name', StringType(),nullable=False)])
BSchema = StructType([StructField('id', IntegerType(),nullable=False),
                 StructField('role', StringType(),nullable=False)])
AData = sc.parallelize ([ Row(1,'michel'), Row(2,'diederik'), Row(3,'rok'), Row(4,'piet')])
BData = sc.parallelize ([ Row(1,'engineer'), Row(2,'lead'), Row(3,'scientist'), Row(5,'manager')])
ADF = hc.createDataFrame(AData,ASchema)
BDF = hc.createDataFrame(BData,BSchema)
DFJOIN = ADF.join(BDF, ADF['id'] == BDF['id'], "outer")
DFJOIN.show()

Input:
+----+--------+----+---------+
|  id|    name|  id|     role|
+----+--------+----+---------+
|   1|  michel|   1| engineer|
|   2|diederik|   2|     lead|
|   3|     rok|   3|scientist|
|   4|    piet|null|     null|
|null|    null|   5|  manager|
+----+--------+----+---------+

from pyspark.sql.functions import *
DFJOINMERGE = DFJOIN.withColumn("_merge", when(ADF["id"].isNull(), "right_only").when(BDF["id"].isNull(), "left_only").otherwise("both"))\
  .withColumn("id", coalesce(ADF["id"], BDF["id"]))\
   .drop(ADF["id"])\
   .drop(BDF["id"])
DFJOINMERGE.show()

Output
+---+--------+---+---------+------+
| id|    name| id|     role|_merge|
+---+--------+---+---------+------+
|  1|  michel|  1| engineer|  both|
|  2|diederik|  2|     lead|  both|
|  3|     rok|  3|scientist|  both|
|  4|    piet|  4|     null|  both|
|  5|    null|  5|  manager|  both|
+---+--------+---+---------+------+

 ==> I would expect id 4 to be left, and id 5 to be right.

Changing join to "left":


Input:
+---+--------+----+---------+
| id|    name|  id|     role|
+---+--------+----+---------+
|  1|  michel|   1| engineer|
|  2|diederik|   2|     lead|
|  3|     rok|   3|scientist|
|  4|    piet|null|     null|
+---+--------+----+---------+

Output
+---+--------+---+---------+------+
| id|    name| id|     role|_merge|
+---+--------+---+---------+------+
|  1|  michel|  1| engineer|  both|
|  2|diederik|  2|     lead|  both|
|  3|     rok|  3|scientist|  both|
|  4|    piet|  4|     null|  both|
+---+--------+---+---------+------+

【问题讨论】:

    标签: python pandas pyspark spark-dataframe


    【解决方案1】:

    试试这个:

    >>> from pyspark.sql.functions import *
    >>> sdf1 = sqlContext.createDataFrame(df1)
    >>> sdf2 = sqlContext.createDataFrame(df2)
    >>> sdf = sdf1.join(sdf2, sdf1["col1"] == sdf2["col1"], "outer")
    >>> sdf.withColumn("_merge", when(sdf1["col1"].isNull(), "right_only").when(sdf2["col1"].isNull(), "left_only").otherwise("both"))\
    ...  .withColumn("col1", coalesce(sdf1["col1"], sdf2["col1"]))\
    ...   .drop(sdf1["col1"])\
    ...   .drop(sdf2["col1"])
    

    【讨论】:

    • 非常感谢。我测试了它,似乎每次都有,我会把测试代码放在问题中。
    • 我无法重现该问题。这两个例子对我来说都很好。
    • 有趣,你用的是哪个 spark 和 python 版本?我们使用 Jupyter notebook 使用 Spark 1.6 和 Python 2,7。如果您运行完全相同的示例代码,则可能是一些版本问题。
    • 1.6.2 / Python 3.5
    • 在 python 3 中得到了相同的结果,但是通过更改答案中的某些字段使其工作。将结果放入问题中
    【解决方案2】:

    改变了 LostInOverflow 的答案并得到了这个工作:

    from pyspark.sql import Row
    
    ASchema = StructType([StructField('ida', IntegerType(),nullable=False),
                     StructField('name', StringType(),nullable=False)])
    BSchema = StructType([StructField('idb', IntegerType(),nullable=False),
                     StructField('role', StringType(),nullable=False)])
    AData = sc.parallelize ([ Row(1,'michel'), Row(2,'diederik'), Row(3,'rok'), Row(4,'piet')])
    BData = sc.parallelize ([ Row(1,'engineer'), Row(2,'lead'), Row(3,'scientist'), Row(5,'manager')])
    ADF = hc.createDataFrame(AData,ASchema)
    BDF = hc.createDataFrame(BData,BSchema)
    DFJOIN = ADF.join(BDF, ADF['ida'] == BDF['idb'], "outer")
    DFJOIN.show()
    
    
    +----+--------+----+---------+
    | ida|    name| idb|     role|
    +----+--------+----+---------+
    |   1|  michel|   1| engineer|
    |   2|diederik|   2|     lead|
    |   3|     rok|   3|scientist|
    |   4|    piet|null|     null|
    |null|    null|   5|  manager|
    +----+--------+----+---------+
    
    from pyspark.sql.functions import *
    DFJOINMERGE = DFJOIN.withColumn("_merge", when(DFJOIN["ida"].isNull(), "right_only").when(DFJOIN["idb"].isNull(), "left_only").otherwise("both"))\
      .withColumn("id", coalesce(ADF["ida"], BDF["idb"]))\
       .drop(DFJOIN["ida"])\
       .drop(DFJOIN["idb"])
    #DFJOINMERGE.show()
    DFJOINMERGE.groupBy("_merge").count().show()
    
    +----------+-----+
    |    _merge|count|
    +----------+-----+
    |right_only|    1|
    | left_only|    1|
    |      both|    3|
    +----------+-----+
    

    【讨论】:

      猜你喜欢
      • 2023-03-10
      • 2020-03-14
      • 2020-09-29
      • 2011-06-26
      • 2018-03-11
      • 2016-09-26
      • 2022-10-13
      • 2022-08-19
      • 2020-12-08
      相关资源
      最近更新 更多