Разбираемся с unionByName() в Spark

Представьте, что у вас есть два датафрейма из разных источников: один выгрузили из PostgreSQL с колонками name, age, city, а второй из CSV с city, name, age. Обычный union() в «лучшем» случае сольёт их по позициям, и city из первого улетит в age второго, либо упадет в ошибку. А unionByName() решает это элегантно: матчит колонки по именам, а не по порядку.

Как работает unionByName()?

unionByName(df2, allowMissingColumns=True) объединяет датафреймы, сопоставляя столбцы по названиям, независимо от их последовательности.

Если в одном датафрейме колонка отсутствует, то подставится null. Это особенно полезно в ETL-пайплайнах, где схемы слегка расходятся (добавились/убрались поля), но семантика та же.

Пример кода df1 = spark.createDataFrame([ ("Anton", 23, "Moscow"), ("Valentina", 27, "Omsk") ], ["name", "age", "city"])

df2 = spark.createDataFrame([ ("Frank", "Moscow", 19, 'M'), ("Boris", "Omsk", 55, 'M') ], ["name", "city", "age", "gender"]) # порядок другой!

result = df1.unionByName(df2, allowMissingColumns=True) result.show()

+---------+---+------+------+ | name|age| city|gender| +---------+---+------+------+ | Anton| 23|Moscow| null| |Valentina| 27| Omsk| null| | Frank| 19|Moscow| M| | Boris| 55| Omsk| M| +---------+---+------+------+

Результат: spark сам разобрался, и все красиво слиплось по именам колонок.

Особенности: ⭐️В случае «плавающих» схем в источниках необходимо добавить параметр allowMissingColumns=True. Тогда при появлении новых полей в одном из датафреймов автоматически добавятся отсутствующие колонки со значениями null для присоединения другого датафрейма. В версиях spark < 3.1.0 придется вручную выравнивать схемы через select или withColumn. ⭐️Добавлять distinct() в конце для union без дублей. Поскольку в spark union() и unionByName() работают как union all (сохраняют дубликаты). ⭐️До соединения датафреймов лучше унифицировать типы данных. Spark будет пытаться привести типы, но может упасть, если это невозможно (например, String vs Int).

Что в итоге?

Считаю, что это один из самых удобных методов в spark. Я практически во всех своих пайплайнах использую unionByName(), чтобы не заморачиваться с дрейфующими схемами данных. Пожалуй, только за исключением данных типа StructType с несколькими уровнями StructType внутри. Как правило, на корневом уровне все срабатывает корректно, но со вложенными структурами уже начинается путаница.

Подробнее в статьях: spark.apache.org, mungingdata.

©️что-то на инженерном


В этом посте были ссылки, но мы их удалили по правилам Сетки