С помощью udf!
🧐 Что это?
UDF (user defined function) - это функция, которую мы сами написали и можем применить в спарке (и не только).
🤩 Пример
Самый обычный джойн:
df1.join(df2, ['id'], 'inner')
План запроса - стандартный SortMerge (без деталей):
== Physical Plan ==
+- Project
+- SortMergeJoin
:- Sort
: +- Exchange hashpartitioning
: +- Filter isnotnull
: +- Scan ExistingRDD
+- Sort ...
😭 Перепишем запрос через udf:
compare_udf = F.udf(lambda x, y: x == y, BooleanType())
df1.alias('df1') \
.join(
df2.alias('df2'),
compare_udf(
F.col('df1.id'),
F.col('df2.id')
),
'inner'
)
И все - теперь у нас под капотом декартово произведение:
== Physical Plan ==
*(3) Project
+- *(3) Filter pythonUDF0#56: boolean
+- BatchEvalPython [<lambda>(id#0L, id#4L)], [pythonUDF0#56]
+- CartesianProduct
:- *(1) Scan ExistingRDD
+- *(2) Scan ExistingRDD
👀 А все потому, что для спарка udf - это черный ящик, и он не будет заглядывать вовнутрь. Так что схема такая:
SortMerge -> Cartesian
ShuffledHash -> Cartesian
BroadcastHash -> BroadcastNestedLoop#spark
