【问题标题】:PySpark: select a column based on the condition another columns values match some specific values, then create the match result as a new columnPySpark:根据条件选择一列,另一列值匹配某些特定值,然后将匹配结果创建为新列
【发布时间】:2020-09-17 02:43:57
【问题描述】:

我之前问过questions的相似度,但由于某些原因,我不得不在PySpark中重新实现它。
例如,

app      col1

app1     anybody love me?
app2     I hate u
app3     this hat is good
app4     I don't like this one
app5     oh my god
app6     damn you.
app7     such nice girl
app8     xxxxx
app9     pretty prefect
app10    don't love me.
app11    xxx anybody?

我想匹配['anybody', 'love', 'you', 'xxx', 'don't']这样的关键字列表,并选择匹配的关键字结果作为新列,命名关键字如下:

app      keyword

app1     [anybody, love]
app4     [don't]
app6     [you]
app8     [xxx]
app10    [don't, love]
app11    [xxx]

作为公认的答案,我可以做的合适方法是创建一个临时数据帧,该数据帧由字符串列表转换,然后 inner join 这两个数据帧一起。
select条件匹配的appkeyword的行。

-- Hiveql implementation
select t.app, k.keyword
from  mytable t
inner join (values ('anybody'), ('you'), ('xxx'), ('don''t')) as k(keyword)
    on t.col1 like conca('%', k.keyword, '%')


但是我不熟悉PySpark 并且很难重新实现它。
谁能帮帮我?
提前致谢。

【问题讨论】:

  • 在 pyspark 中,您可以轻松编写一个获取 col1 字符串的 UDF,然后它只是简单的 python 代码(拆分字符串)并做您需要的事情。

标签: string apache-spark pyspark apache-spark-sql hiveql


【解决方案1】:

请在下面找到两种可能的方法:

选项 1

第一个选项是使用数据框 API 来实现与上一个问题中类似的连接。这里我们将keywords 列表转换成一个数据帧,然后将它与大数据帧连接起来(注意我们广播了小数据帧以确保更好的性能):

from pyspark.sql.functions import broadcast

df = spark.createDataFrame([
  ["app1", "anybody love me?"],
  ["app4", "I don't like this one"],
  ["app5", "oh my god"],
  ["app6", "damn you."],
  ["app7", "such nice girl"],
  ["app8", "xxxxx"],
  ["app10", "don't love me."]
]).toDF("app", "col1")

# create keywords dataframe
kdf = spark.createDataFrame([(k,) for k in keywords], "key string")

# +-----+
# |  key|
# +-----+
# |  xxx|
# |don't|
# +-----+

df.join(broadcast(kdf), df["col1"].contains(kdf["key"]), "inner")

# +-----+---------------------+-----+
# |app  |col1                 |key  |
# +-----+---------------------+-----+
# |app4 |I don't like this one|don't|
# |app8 |xxxxx                |xxx  |
# |app10|don't love me.       |don't|
# +-----+---------------------+-----+

连接条件基于Column类的contains函数。

选项 2

您还可以在 expr 中将 PySpark 高阶函数 filterrlike 结合使用:

from pyspark.sql.functions import lit, expr, array

df = spark.createDataFrame([
  ["app1", "anybody love me?"],
  ["app4", "I don't like this one"],
  ["app5", "oh my god"],
  ["app6", "damn you."],
  ["app7", "such nice girl"],
  ["app8", "xxxxx"],
  ["app10", "don't love me."]
]).toDF("app", "col1")

keywords = ["xxx", "don't"]

df.withColumn("keywords", array([lit(k) for k in keywords])) \
  .withColumn("keywords", expr("filter(keywords, k -> col1 rlike k)")) \
  .where("size(keywords) > 0") \
  .show(10, False)

# +-----+---------------------+--------+
# |app  |col1                 |keywords|
# +-----+---------------------+--------+
# |app4 |I don't like this one|[don't] |
# |app8 |xxxxx                |[xxx]   |
# |app10|don't love me.       |[don't] |
# +-----+---------------------+--------+

解释

  1. 我们使用array([lit(k) for k in keywords]) 生成一个数组,其中包含我们的搜索所基于的关键字,然后我们使用withColumn 将其附加到现有数据框。

  2. 接下来是expr("size(filter(keywords, k -> col1 rlike k)) > 0"),我们遍历关键字项,试图找出其中是否存在于 col1 文本中。如果这是真的,filter 将返回一个或多个项目,size 将大于 0,这构成了我们用于检索记录的 where 条件。

【讨论】:

  • @Bowen 完全没有问题,你能详细说明一下吗?什么是预期的输出?例如给定选项 2,这是预期的输出吗?
  • 所以给定关键字 ["xxx", "don't"] 我们在 col1 中搜索包含它们的文本。在上面的示例中,app4 和 app8 满足该条件。
  • @Bowen 现在您正在调用df.show(),它正在打印原始数据帧。由于 RDD 是不可变的,因此连接转换不会修改原始 df。为了得到你需要做的最终结果final_df = df.join(broadcast(kdf), df["col1"].contains(kdf["key"]), "inner") 然后final_df.show()
  • 对于选项 2,错误 here 表示未找到名称为 col 的列,可能您输入错误,因为没有这样的列。尝试再次复制粘贴代码
  • @Bowen 选项 2 确实会将所有匹配项存储到 keywords col
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-03-04
  • 1970-01-01
  • 2021-03-28
  • 1970-01-01
  • 2021-08-16
  • 1970-01-01
相关资源
最近更新 更多