【问题标题】:Compare String Values in 2 Spark Dataframes [duplicate]比较 2 个 Spark 数据帧中的字符串值 [重复]
【发布时间】:2018-05-23 07:49:14
【问题描述】:

我有 2 个名为 -brand_name 和 poi_name 的数据框。

数据框 1(品牌名称):-

+-------------+
|brand_stop[0]|
+-------------+
|TOASTMASTERS |
|USBORNE      |
|ARBONNE      |
|USBORNE      |
|ARBONNE      |
|ACADEMY      |
|ARBONNE      |
|USBORNE      |
|USBORNE      |
|PILLAR       |
+-------------+

数据框 2:-(poi_name)

+---------------------------------------+
|Name                                   |
+---------------------------------------+
|TOASTMASTERS DISTRICT 48               |
|USBORNE BOOKS AND MORE                 |
|ARBONNE                                |
|USBORNE BOOKS AT HOME                  |
|ARBONNE                                |
|ACADEMY, LTD.                          |
|ARBONNE                                |
|USBORNE BOOKS AT HOME                  |
|USBORNE BOOKS & MORE                   |
|PILLAR TO POST HOME INSPECTION SERVICES|
+---------------------------------------+

我想检查数据框 1 的 brand_stop 列中的字符串是否存在于数据框 2 的名称列中。匹配应该逐行完成,然后如果匹配成功,则该特定记录应存储在新的柱子。

我尝试使用 Join 过滤数据框:-

from pyspark.sql.functions import udf, col 
from pyspark.sql.types import BooleanType

contains = udf(lambda s, q: q in s, BooleanType())

like_with_python_udf = (poi_names.join(brand_names1)
    .where(contains(col("Name"), col("brand_stop[0]")))
    .select(col("Name")))
like_with_python_udf.show()

但这显示错误

"AnalysisException: u'Detected cartesian product for INNER join between logical plans"

我是 PySpark 的新手。请帮我解决这个问题。

谢谢

【问题讨论】:

  • 加入应该为您解决问题
  • 区分大小写怎么样?
  • @Steven 考虑到数据框 1 元素是大写的。之后请建议算法。
  • @RameshMaharjan 无法在这两个数据帧之间加入,因为这两个数据帧中没有“id”。如果有其他方式加入,请告诉我。
  • 在上面的示例中,输出将与 Dataframe 2 相同,因为所有行都匹配成功。但是在我正在处理的原始数据框中(我的意思是两个数据框都有 6000 行 - 可能有一些不匹配的行)。在上面的例子中,我只是展示了两个数据帧的前 10 行,只是为了说明这个例子。因此,在输出数据帧中 - 只有两个数据帧之间的匹配行应该存在。

标签: python apache-spark dataframe pyspark apache-spark-sql


【解决方案1】:

scala 代码会是这样的:

val d1 = Array(("TOASTMASTERS"),("USBORNE"),("ARBONNE"),("USBORNE"),("ARBONNE"),("ACADEMY"),("ARBONNE"),("USBORNE"),("USBORNE"),("PILLAR"))
val rdd1 = sc.parallelize(d1)
val df1 = rdd1.toDF("brand_stop")

val d2 = Array(("TOASTMASTERS DISTRICT 48"),("USBORNE BOOKS AND MORE"),("ARBONNE"),("USBORNE BOOKS AT HOME"),("ARBONNE"),("ACADEMY, LTD."),("ARBONNE"),("USBORNE BOOKS AT HOME"),("USBORNE BOOKS & MORE"),("PILLAR TO POST HOME INSPECTION SERVICES")) 
val rdd2 =sc.parallelize(d2)
val df2 = rdd2.toDF("names")


def matchFunc(s1:String,s2:String) : Boolean ={ 
if(s2.contains(s1)) true
else false
}
val contains = udf(matchFunc _)

val like_with_python_udf = (df1.join(df2).where(contains(col("brand_stop"), col("names"))).select(col("brand_stop"), col("names")))
like_with_python_udf.show()

Python 代码:

from pyspark.sql import Row
from pyspark.sql.functions import udf, col 
from pyspark.sql.types import BooleanType

schema1 = Row("brand_stop")
schema2 = Row("names")

df1 = sc.parallelize([
    schema1("TOASTMASTERS"),
    schema1("USBORNE"),
    schema1("ARBONNE")
]).toDF()
df2 = sc.parallelize([
    schema2("TOASTMASTERS DISTRICT 48"),
    schema2("USBORNE BOOKS AND MORE"),
    schema2("ARBONNE"),
    schema2("ACADEMY, LTD."),
    schema2("PILLAR TO POST HOME INSPECTION SERVICES")
]).toDF()

contains = udf(lambda s, q: q in s, BooleanType())

like_with_python_udf = (df1.join(df2)
    .where(contains(col("brand_stop"), col("names")))
    .select(col("brand_stop"), col("names")))
like_with_python_udf.show()

我得到输出:

+------------+ |品牌站| +------------+ |祝酒师| | USBORNE| |阿邦| +------------+

【讨论】:

  • 感谢@Mugdha 的回复。当我运行它时(将其转换为 Python 后),显示的错误是:-“AnalysisException: u'Detected cartesian product for INNER join between logical plans”。请帮我解决这个问题。
  • 嘿@AnubhavSarangi,我已经用python代码编辑了答案。可以查一下吗?
【解决方案2】:

匹配应该逐行进行

在这种情况下,您必须添加某种形式的索引并加入

from pyspark.sql.types import *

def index(df):
    schema = StructType(df.schema.fields + [(StructField("_idx", LongType()))])
    rdd = df.rdd.zipWithIndex().map(lambda x: x[0] +(x[1], ))
    return rdd.toDF(schema)

brand_name = spark.createDataFrame(["TOASTMASTERS", "USBORNE"], "string").toDF("brand_stop")
poi_name = spark.createDataFrame(["TOASTMASTERS DISTRICT 48", "USBORNE BOOKS AND MORE"], "string").toDF("poi_name")

index(brand_name).join(index(poi_name), ["_idx"]).selectExpr("*", "poi_name rlike brand_stop").show()
# +----+------------+--------------------+-------------------------+              
# |_idx|  brand_stop|            poi_name|poi_name RLIKE brand_stop|
# +----+------------+--------------------+-------------------------+
# |   0|TOASTMASTERS|TOASTMASTERS DIST...|                     true|
# |   1|     USBORNE|USBORNE BOOKS AND...|                     true|
# +----+------------+--------------------+-------------------------+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-08-14
    • 1970-01-01
    • 2019-09-08
    • 2023-04-01
    • 1970-01-01
    • 2018-01-15
    • 1970-01-01
    • 2023-03-26
    相关资源
    最近更新 更多