【问题标题】:PySpark: Compare two Dataframes based on common String column and generate Result Boolean withColumn()PySpark:基于公共字符串列比较两个数据框并生成结果布尔 withColumn()
【发布时间】:2020-02-04 16:08:42
【问题描述】:

我有两个数据帧 ddf_1 和 ddf_2 共享一个字符串 ID 的公共列。我的目标是在 ddf_1 中创建一个新的布尔 is_fine 列,如果 ID 包含在 ddf_1 和 ddf_2 中,则该列包含 True,如果 ID 不包含在 ddf_1 和 ddf_2 中,则包含 False。

考虑这个示例数据:

#### test
#example data
data_1 = { 
    'fruits': ["apples", "banana", "cherry"],
    'myid': ['1-12', '2-12', '3-13'],
    'meat': ["pig", "cow", "chicken"]}

data_2 = { 
    'furniture': ["table", "chair", "lamp"],
    'myid': ['1-12', '0-11', '2-12'],
    'clothing': ["pants", "shoes", "socks"]}

df_1 = pd.DataFrame(data_1)
ddf_1 = spark.createDataFrame(df_1)

df_2 = pd.DataFrame(data_2)
ddf_2 = spark.createDataFrame(df_2)

我想象一个类似这样的函数:

def func(df_1, df_2, column_1, column_2):
    if df_1.column_1 != df_2.column_2:
       return df_1.withColumn('is_fact', False)
    else:
        return df_1.withColumn('is_fact', True)
    return df_1

所需的输出应如下所示:

【问题讨论】:

标签: python join pyspark comparison conditional-statements


【解决方案1】:

您可以在my_id 列上的2 个数据帧之间执行左外连接,并使用简单的case 语句推导出is_fine 列,如下所示,

import pyspark.sql.functions as F

ddf_1.join(ddf_2, ddf_1.myid == ddf_2.myid, 'left')\
.withColumn('is_fine', F.when(ddf_2.myid.isNull(), False).otherwise(True))\
.select(ddf_1['fruits'], ddf_1['myid'], ddf_1['meat'], 'is_fine').show()

输出:

+------+----+-------+-------+
|fruits|myid|   meat|is_fine|
+------+----+-------+-------+
|cherry|3-13|chicken|  false|
|apples|1-12|    pig|   true|
|banana|2-12|    cow|   true|
+------+----+-------+-------+

【讨论】:

  • 谢谢@noufel13 这对我很有用。我唯一改变的是使列选择语句更通用,因为这是假数据,我通常有很多列要选择,并且不能单独对每个列进行硬编码。您可以在下面找到代码
【解决方案2】:

利用 Spark SQL 解决此类问题:

query = """
select ddf_1.*,
case 
    when ddf_1.myid = ddf_2.myid  then True
    else False 
end as is_fine
from ddf_1 left outer join ddf_2 
on ddf_1.myid = ddf_2.myid
"""

display(spark.sql(query))

这是output

【讨论】:

  • 我会做 'CASE WHEN ddf_2.myid IS NULL' 而不是 'when ddf_1.myid = ddf_2.myid'
  • 我得到一个表或视图未找到:ddf_1;第 8 行 pos 5 异常。
  • 您需要在ddf_1.createOrReplaceTempView("ddf_1")之前将数据框转换为临时视图
【解决方案3】:
#left join ddf2 on ddf1
result = (ddf_1.join(ddf_2, ddf_1.myid == ddf_2.myid, how='left')\
          #create is_fine column
          .withColumn('is_fine', F.when(ddf_2.myid.isNull(), False).otherwise(True)))\
          #select all columns from ddf_1, the new column is_fine and show
          .select(ddf_1["*"], "is_fine").show() 

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-02
    • 2021-01-16
    • 1970-01-01
    相关资源
    最近更新 更多