Python UDF в PySpark: почему дорого и чем заменить
Обычный Python UDF дорог сразу по двум причинам. Механическая: данные из JVM сериализуются, передаются в отдельный Python-процесс построчно, обрабатываются и возвращаются обратно — накладные расходы на каждую строку. Логическая, и она хуже: для Catalyst UDF непрозрачен, поэтому проталкивание предикатов и отсечение колонок вокруг него перестают работать — фильтр по UDF заставляет прочитать весь датасет. Порядок замены такой. Первым делом ищите встроенную функцию из pyspark.sql.functions — их сотни, и они выполняются внутри JVM с генерацией кода. Не нашли — берите pandas_udf: данные передаются пакетами через Arrow, и это обычно в разы быстрее построчного UDF. В Spark 4 к ним добавились Arrow-оптимизированные Python UDF, работающие напрямую с PyArrow без прослойки pandas. И только если ничего не подошло, пишите обычный UDF — но никогда не ставьте его в фильтр.