IIUC, you want to return the rows in which column_a is "like" (in the SQL sense) any of the values in list_a.
One way is to use functools.reduce:
from functools import reduce
list_a = ['string', 'third']
df1 = df.where(
reduce(lambda a, b: a|b, (df['column_a'].like('%'+pat+"%") for pat in list_a))
)
df1.show()
#+------------+-----+
#| column_a|count|
#+------------+-----+
#| some_string| 10|
#|third_string| 30|
#+------------+-----+
Essentially you loop over all of the possible strings in list_a to compare in like and "OR" the results. Here is the execution plan:
df1.explain()
#== Physical Plan ==
#*(1) Filter (Contains(column_a#0, string) || Contains(column_a#0, third))
#+- Scan ExistingRDD[column_a#0,count#1]
Another option is to use pyspark.sql.Column.rlike instead of like.
df2 = df.where(
df['column_a'].rlike("|".join(["(" + pat + ")" for pat in list_a]))
)
df2.show()
#+------------+-----+
#| column_a|count|
#+------------+-----+
#| some_string| 10|
#|third_string| 30|
#+------------+-----+
Which has the corresponding execution plan:
df2.explain()
#== Physical Plan ==
#*(1) Filter (isnotnull(column_a#0) && column_a#0 RLIKE (string)|(third))
#+- Scan ExistingRDD[column_a#0,count#1]
Answer from pault on Stack Overflowwhat it says is "df.score in l" can not be evaluated because df.score gives you a column and "in" is not defined on that column type use "isin"
The code should be like this:
# define a dataframe
rdd = sc.parallelize([(0,1), (0,1), (0,2), (1,2), (1,10), (1,20), (3,18), (3,18), (3,18)])
df = sqlContext.createDataFrame(rdd, ["id", "score"])
# define a list of scores
l = [10,18,20]
# filter out records by scores by list l
records = df.filter(~df.score.isin(l))
# expected: (0,1), (0,1), (0,2), (1,2)
# include only records with these scores in list l
df.filter(df.score.isin(l))
# expected: (1,10), (1,20), (3,18), (3,18), (3,18)
Note that where() is an alias for filter(), so both are interchangeable.
based on @user3133475 answer, it is also possible to call the isin() function from col() like this:
from pyspark.sql.functions import col
l = [10,18,20]
df.filter(col("score").isin(l))
To make your function work , you should create an array column to compare:
df.select(fn.array([fn.lit(i) for i in key_labels])).show(truncate=False)
+----------------------------------+
|array(COMMISSION, COM, PRET, LOAN)|
+----------------------------------+
|[COMMISSION, COM, PRET, LOAN] |
|[COMMISSION, COM, PRET, LOAN] |
+----------------------------------+
So you code would look like below:
def containsAny(string, array):
if len(string) == 0:
return False
else:
return (any(word in string for word in array))
contains_udf = fn.udf(containsAny, T.BooleanType())
(df.withColumn("keyword_match", contains_udf(fn.col("original"),
fn.array([fn.lit(i) for i in key_labels])))).show()
Outputs:
+----------+---+-------------+
| original| id|keyword_match|
+----------+---+-------------+
|COMMISSION| 1| true|
|CAMMISSION| 2| false|
+----------+---+-------------+
However you could also use isin:
df.withColumn('keyword_match',df['original'].isin(key_labels)).show()
+----------+---+-------------+
| original| id|keyword_match|
+----------+---+-------------+
|COMMISSION| 1| true|
|CAMMISSION| 2| false|
+----------+---+-------------+
Another solution that works as well is rlike function. In fact it works much faster than an udf.
regex = "|".join(r"(" + x + r")" for x in key_labels)
df = spark.createDataFrame([("COMMISSION", "1"), ("CAMMISSION", "2")], ("original", "id"))
df.select("original","id",fn.col("original").rlike(regex).alias("keyword_match")).show()