Привет, Хабр! Я Павел Молчанов, Data Engineer из группы обработки гео- и медиаданных Big Data компании МТС Web Services. Наш продукт «ГеоЭффект» занимается расчетами различных геостатистик на основе телеком-данных (например, по туризму).
У нас есть ежедневно обновляемая инфраструктурной командой витрина-источник с событиями (точками нахождения абонентов) на основе мобильной сети по всем регионам присутствия МТС. За сутки на каждого абонента приходится в среднем порядка 200 геособытий, общий объем данных в витрине — около 11 млрд записей в сутки. Однако в витрине-источнике не было привязки геособытий к административно-территориальной информации «Регион РФ / район», что критически важно нам как продуктовой команде, решающей бизнес-кейсы в интересах заказчиков.
Отсюда и появилась задача — дообогатить витрину дополнительными атрибутами на основе координат геособытий. Так что в этом материале по мотивам доклада для Highload++ поделюсь, как мы это сделали, а еще ускорили джойн датафреймов с геоданными.
Пара слов по технологическому стеку: сейчас мы активно проводим подготовительные работы по переходу на Data LakeHouse, но пока что основная рабочая лошадка — Spark на старом добром Hadoop-кластере с HDFS и Hive Metastore. Так что повествование будет идти в контексте технологий Spark on YARN + Sedona + HDFS + Hive.

Что за Apache Sedona
Sedona ранее была известна как GeoSpark и позднее переименована в Apache Sedona. По сути, это расширение для Apache Spark для распределенной обработки пространственных геоданных. По поводу собственно Apache Spark: почти все команды в Big Data MTS работают на PySpark. Scala Spark используется, но эпизодически.
Чтобы подружить Sedona с PySpark, нужно установить в используемое вами виртуальное окружение Python (venv) Python-обертку для Sedona из PyPI. В контексте исследований и бенчмарков, которые будут дальше по ходу повествования, у нас использовался Spark 3.3.1 и Sedona 1.7.2.
В конфиг Spark-сессии нужно добавить парочку джарников. Эти JAR, по сути — библиотека с ядром Sedona и обертки для геометрических типов данных, которые Sedona умеет обрабатывать.
"spark.jars.packages": [ "org.apache.sedona:sedona-spark-shaded-3.3_2.12:1.7.2", "org.datasyslab:geotools-wrapper:1.5.1-28.2" ] … "spark.serializer": "org.apache.spark.serializer.KryoSerializer" "spark.kryo.registrator": "org.apache.sedona.core.serde.SedonaKryoRegistrator"
После старта Spark нужно зарегистрировать Sedona в рамках Spark-сессии — тогда становятся доступными специфические функции для обработки геоданных, о них мы поговорим чуть позже.
from sedona.spark import SedonaContext sedona = SedonaContext.create(spark)
Форматы данных, с которыми мы постоянно работаем и которые важны в контексте рассказа, — это географические координаты, долгота и широта. Также далее будем говорить и о специфических форматах геометрий:
-
WKT (Well-known text). Текстовый формат представления векторных геометрий. Речь идет о некоторых геометрических понятиях:
точка на карте, имеющая формат Point в WKT;
ломаная линия LineString (с ее помощью могут отображаться, например, дороги);
Polygon — замкнутая область на карте.
Простые типы: Point, LineString, Polygon.
Мультигеометрии: MultiPolygon (строго типизированный массив из полигонов) и GeometryCollection (список, который может содержать разные типы геометрий).

Как мы работали с данными раньше
Формат нашей упомянутой во вступлении витрины с геособытиями выглядит так:

Здесь есть идентификатор абонента, дата и время, долгота, широта и добавлена точка POINT уже в формате WKT.
Второй датафрейм — это справочник с прямоугольными полигонами, на которые разбита вся территория присутствия МТС. Мы берем территорию РФ и каждый регион разбиваем в данном контексте на прямоугольную сетку.

На рисунке — «прямоугольная сетка полигонов 500 × 500 м» из справочника. Из нюансов: на долготе/широте Москвы прямоугольные полигоны вырождены в квадратные. Анклав с желтыми квадратиками — например, территория Сколково, которая административно относится к Москве.
Справочник таких прямоугольных полигонов у нас в Big Data построен методом QUADTREE, то есть вложением друг в друга прямоугольников разных размерностей. В прямоугольник 500 на 500 в Москве вложено четыре прямоугольника 250 на 250 и так далее.
У каждого из полигонов есть атрибуты «ID полигона», «регион», «район» и «геометрия».

Что решили улучшить
Мы решили дообогатить геособытия по абонентам из 1-й витрины атрибутами «регион РФ», «район» из 2-го справочника полигонов.
Условие — вхождение точки POINT геособытия из 1-й витрины в полигон POLYGON из 2-го справочника. Все достаточно тривиально: в какой из прямоугольных полигонов войдет точка — к тому региону и району она привязана. Такой join мы часто делаем по условию ST_Contains из API Apache Sedona. По сути, это предикат, который возвращает булево значение True или False. Рассмотрим документацию по функции-предикату ST_Contains с официального сайта Apache Sedona. Здесь А — это левая геометрия, которой является полигон, B — правая геометрия, которой является точка. Если геометрия A полигона содержит точку, то возвращается True. Все просто.

В разделе с предикатами для Apache Sedona есть и другие типы логических условий. Например, ST_Intersects — это функция-предикат, которая возвращает True, если два полигона пересекаются друг с другом, и так далее.
Ниже приведена еще одна визуализация того, чего мы добиваемся в результате джойна:

На карте отрисованы прямоугольные полигоны по Москве, их раскраска (на картинке — желтая или коричневая) привязана к ID полигона. В виде точек здесь отображены некоторые геособытия. Точки попадают в им соответствующие полигоны и привязываются к региону/району.
Желаемый формат данных в виде таблицы выглядит так:

Проблема медленных джойнов дата-фреймов Spark с геоданными на Sedona
Типовая задача по обработке (в том числе join’у) Spark DataFrame’ов с геоданными у нас часто решается на следующих ресурсах:
from sedona.utils import SedonaKryoRegistrator, KryoSerializer "spark.executor.cores": "4", "spark.executor.memory": "12g", "spark.dynamicAllocation.enabled": "false", "spark.executor.instances": "200", "spark.sql.shuffle.partitions": "800", "spark.default.parallelism": "800", "spark.jars.packages": [ "org.apache.sedona:sedona-spark-shaded-3.3_2.12:1.7.2", "org.datasyslab:geotools-wrapper:1.5.1-28.2", ], "spark.serializer": KryoSerializer.getName, "spark.kryo.registrator": SedonaKryoRegistrator.getName,
На Hadoop-кластере у нас доступны в сумме не менее 2,5 ТБ оперативной памяти и не менее 800 ядер — солидные ресурсы на задачку с джойном. Для бенчмарка взята статическая аллокация ресурсов, то есть «прибиты гвоздями» 200 экзекьюторов, на каждом по 4 ядра и по 12 ГБ оперативной памяти. Добавлены джарники Sedona.

Так выглядит простой код в стиле Spark DataFrame API по джойну витрины с геособытиями и справочника с полигонами:
from sedona.sql.st_predicates import ST_Contains # sdf_events - 11_770_557_003 rows # sdf_polygons - 72_650_028 rows sdf_joined = ( sdf_polygons.alias("p") .join( sdf_events.alias("e"), ST_Contains(col("p.geometry"), col("e.point")), "inner" ) .select( col("e.hash"), col("e.time"), col("e.point").cast("string").alias("point"), col("p.region_name"), col("p.district_name") ) )
sdf_polygons — это 72 млн прямоугольных полигонов из справочника, а sdf_events — это более 11 млрд геособытий по всей абонентской базе МТС за один день. Join осуществляется по условию ST_Contains, после чего в итоговую выборку включаются поля из обоих датафреймов.
Заглянув в Spark User Interface, мы увидим очень неприятную картину: расчет продолжается уже почти 1,5 часа, висят 14 последних тасок.

Если заглянем в метрики Spark Tasks в табличку с перцентилями, то эти опасения подтверждаются: 75-й перцентиль Spark-задач выполняется за 1,5 минуты, тогда как задачи из 100-го перцентиля зависают почти на 1,5 часа. Распределение времени выполнения по исполнителям явно указывает на сильный перекос данных (Data Skew).

С этой проблемой при решении задач на нашем продукте сначала сталкивались дата-аналитики, затем к ее решению подключились дата-инженеры. Было понятно, что причина в перекосе данных, но универсального способа борьбы с ним не находилось. Я стал копать глубже.
В нашем случае используются пространственные предикаты вроде ST_Contains и ST_Intersects. Классические методы борьбы с Data Skew применить нельзя — Join выполняется не по равенству ключей, а по геометрическому условию.
Готового решения в открытых источниках я не нашел, но выделил два перспективных направления:
на Stack Overflow узнал, что Apache Sedona поддерживает не только DataFrame API, но и более низкоуровневую абстракцию SpatialRDD. Она является частью ядра Sedona и позволяет использовать пространственное индексирование при выполнении джойнов.
Оба подхода показались перспективными векторами действий по фиксу проблемы медленных джоинов на Sedona, поэтому я начал более глубокое их исследование. За день построил пространственные сетки на данных МТС по Москве и Московской области.
В процессе оказалось, что в библиотеке Sedona «из коробки» можно построить всего три типа пространственных сеток:
1. EQUALGRID (равномерная сетка) — простой вариант. Исходим из того, что задаем желаемое количество ячеек сетки разбиения, Sedona по BoundingBox вокруг рассматриваемой области строит прямоугольную сетку с одинаковой площадью ячейки.

2. QUADTREE — вариант поинтереснее. Берем датафрейм с геособытиями по Москве и области и строим сетку типа квадродерева. В контексте: не следует путать ее с упоминавшейся выше сеткой из справочника. В последней может быть что угодно — гексагоны или произвольная геометрия с жилой застройкой, лесами и так далее. Ниже же служебная сетка, которая строится внутри Sedona, однако ее можно вытащить и визуализировать, что мы и сделали ?

Что здесь интересного? Местность разбивается на квадраты либо прямоугольники, причем есть вложенность. В квадрат большого размера вкладываются четыре маленьких, то есть каждый квадрат может быть рекурсивно разбит на четыре квадрата меньшего размера. Sedona обучилась на датафрейме с геособытиями по Москве и области, и в центре Москвы — малюсенькие квадратики, по краям — геометрические области, огромные по площади.
На этом моменте, после построения пространственных сеток разных типов, я осознал причину всех наших сложностей: с одной стороны, Sedona умеет строить сетку по равномерному распределению площади, с другой стороны, Sedona умеет делать это по примерно равномерному распределению событий по ячейкам. Далее ячейки раскидываются по экзекьюторам и по Spark-таскам, одна ячейка обрабатывается на отдельном ядре в виде Spark-партиции. Собственно, из-за неправильного построения пространственной сетки у нас и возникал перекос данных.
3. KDBTREE — это k-мерное бинарное дерево, которое также строится на плоскости. Здесь даже невооруженным глазом видно, что Bounding Box был сначала поделен по горизонтали, потом по вертикали, затем прямоугольники делятся пополам. Принцип тот же: события раскидываются примерно равномерно по ячейкам.

Итак, я осознал причину всех наших предыдущих бедствий с Data Skew. Если визуализировать наш датафрейм с 11 млрд геособытий в виде тепловой карты по привязке к местности (к регионам РФ, в которых присутствует компания МТС), то получается следующее: при неравномерном распределении событий на местности причина Data Skew при джойне кроется в неправильно построенной сетке.

Очевидно, геособытия сконцентрированы там, где живет большая часть населения — то есть к западу от Уральских гор и по южной границе России. Невооруженным глазом видны яркие выбросы в мегаполисах (Москва, Санкт-Петербург).
Таким образом, причина наших проблем оказалась в том, что под капотом Sedona по умолчанию стоила что-то вроде равномерной сетки, обучаясь по справочнику равномерных прямоугольных полигонов, а не по датафрейму, в котором было много неравномерно распределенных геособытий.

Секрет был раскрыт, и далее я приступил к написанию нашего кастомного инструмента по быстрым джойнам на Sedona.
Разработка инструмента по быстрым индексированным джойнам на Sedona
Изначально я вдохновился приведенными на Stack Overflow примерами работы со SpatialRDD и решил написать универсальный инструмент, который могли бы использовать наши коллеги во всей геокоманде МТС.
Работу над кастомным инструментом быстрых индексированных Spatial Join построил следующим образом (ниже используется Scala):
Преобразовал оба Spark DataFrame в SpatialRDD через Adapter.toSpatialRdd и выполнил analyze().
Сохранил все негеометрические поля (leftSdf.columns/rightSdf.columns), иначе после join остались бы только геометрии.
Выбрал сторону построения сетки (JoinBuildSide). Для неравномерно распределенных данных сетку обычно строят на большем датафрейме.
Выполнил spatialPartitioning(), задав тип сетки (QUADTREE, KDBTREE и другие) и количество партиций numGridPartitions.
Применил построенный партиционер ко второму SpatialRDD. Оба RDD должны использовать одну и ту же сетку.
Построил пространственные индексы через buildIndex().
Выполнил SpatialJoinQueryFlat с нужным пространственным предикатом (ST_Contains, ST_Intersects, ST_Touches и другие).
Извлек атрибуты из UserData, так как результат джойна хранится в Geometry и без дополнительной обработки не содержит исходных колонок.
Преобразовал результат обратно в DataFrame через toDF().
После возврата в DataFrame восстановил исходные поля из массивов rightData и leftData через select и cast.
Ключевые параметры функции
leftSdf/rightSdf — входные датафреймы; leftGeomField/rightGeomField — геометрические поля; numGridPartitions — число ячеек сетки; gridType — тип пространственной сетки; indexType — тип индекса; spatialPredicate — пространственное условие; joinBuildSide — сторона построения сетки.
Основной код
// Scala Code import org.apache.sedona.core.enums. {GridType, IndexType, JoinBuildSide} import org.apache.sedona.core.spatialOperator.{JoinQuery, SpatialPredicate} import org.apache.sedona.sql.utils.Adapter import org.apache.spark.sql.{DataFrame, SparkSession, functions => sf} // Starting Spark Job - aggregate at SpatialRDD val rddLeft = Adapter.toSpatialRdd(leftSdf, leftGeomField, leftSdf.columns) rddLeft.analyze() // Starting Spark Job - aggregate at SpatialRDD val rddRight = Adapter.toSpatialRdd(rightSdf, rightGeomField, rightSdf.columns) rddRight.analyze() joinBuildSide match { case JoinBuildSide.LEFT => { rddLeft.spatialPartitioning(gridType, numGridPartitions) rddRight.spatialPartitioning(rddLeft.getPartitioner) } case JoinBuildSide.RIGHT => { rddRight.spatialPartitioning(gridType, numGridPartitions) rddLeft.spatialPartitioning(rddRight.getPartitioner) } } // Indexing rddLeft.buildIndex(indexType, true) rddRight.buildIndex(indexType, true) val spatialJoinResult = JoinQuery.SpatialJoinQueryFlat(rddLeft, rddRight, true, spatialPredicate) val rddResult = spatialJoinResult.rdd.map{ case (x, y) => ( x.toString, x.getUserData.toString.split("\t"), y.toString, y.getUserData.toString.split("\t") ) } val sdfResult = rddResult.toDF("rightGeom", "rightData", "leftGeom", "leftData")
Пример вызова:
val joinedDf = { joinDfsAsSpatialRdds( leftSdf = sdfPolygons, leftGeomField = "geometry", rightSdf = sdfEvents, rightGeomField = "point", numGridPartitions = 12800, gridType = GridType.QUADTREE, indexType = IndexType.QUADTREE, spatialPredicate = SpatialPredicate.CONTAINS, joinBuildSide = JoinBuildSide.RIGHT ) .select( sf.col("rightData").getItem(0).cast("bigint").alias("hash"), sf.col("rightData").getItem(1).cast("timestamp").alias("time"), sf.col("rightGeom").alias("point"), sf.col("leftData").getItem(1).alias("region_name"), sf.col("leftData").getItem(2).alias("region_code") ) }
В результате при корректном подборе параметров пространственной сетки время выполнения spatial join сократилось примерно с 1,5 часа до 15 минут на данных с выраженным data skew.
Казалось бы, можно наливать шампанского в бокалы, но не тут-то было. Я пришел к своей команде, показываю эту функцию, а ребята говорят: «Да, это работает, но слишком сложно. Не мог бы ты адаптировать это в более user-friendly-формат, чтобы можно было использовать на PySpark с менее высоким входным порогом?» Так что пришлось пойти думать дальше.
Второй подход
В блоке кода выше с примером вызова кастомной функции имеется ряд параметров — расскажу об их подборе. Например, что означает numGridPartitions в вызове функции? Почему он 12 800?

Здесь крайне важно подобрать правильное число ячеек в сетке разбиения. Оценка следующая.
Во-первых, у нас было 800 ядер в расчете. Соответственно, единовременно могло выполняться 800 Spark-тасок, то есть обрабатываться 800 партиций Spark-датафрейма. 12 800 — это 800 * 24 = 16. То есть на каждое ядро у нас приходится по 16 тасок.
Во-вторых, площадь России порядка 17 млн кв. км. Если разделить площадь России на количество ячеек в сетке разбиения, полагая, что разбиение равномерное (если использовать сетку EQUALGRID), получим площадь одной ячейки 1328 кв. км. Как мы помним, все проблемы и беды возникали потому, что мегаполисы попадали в одну или несколько Spark-тасок. Здесь при таком уровне параллелизма площадь ячейки при равномерном разбиении сопоставима с площадью Москвы.
Очевидно, что используя QUADTREE или k-мерное бинарное дерево, мы можем решить проблему перекоса данных при таком уровне ячеек в сетке разбиения, потому что часть ячеек, например по Якутии, будет с огромной площадью, а по Москве и Петербургу — с маленькой размерностью, но с большим количеством событий.
Если подобрать слишком низкий параметр numGridPartitions, у нас будет выраженный случай с перекосом данных, который ранее происходил по умолчанию. Тогда будет слишком большая площадь ячеек и мегаполисы будут попадать снова на всего на пару-тройку ядер.
Как выбирать gridType, indexType и spatialPredicate?
spatialPredicate — один из видов логических джойнов из библиотеки Apache Sedona, которые зашиты в Java API Sedona. Я в основном работал с тремя типами предикатов: это CONTAINS, INTERSECTS и TOUCHES.
org.apache.sedona.core.spatialOperator.SpatialPredicate
CONTAINSINTERSECTSWITHINCOVERSCOVERED_BYTOUCHESOVERLAPSCROSSESEQUALS
gridType — это один из трех типов сетки, который умеет строить Sedona (см. визуализации выше). Очевидно, что в рамках нашей задачи EQUALGRID совершенно не подходит, остается либо квадратное дерево, либо k-мерное бинарное дерево. Я сосредоточился на QUADTREE.
org.apache.sedona.core.enums.GridType
EQUALGRIDQUADTREEKDBTREE
В indexType доступны два вида — QUADTREE и RTREE. В бенчмарках и на проде теперь использую QUADTREE и для сетки, и для индекса.
org.apache.sedona.core.enums.IndexType
QUADTREERTREE
Что в итоге сработало
Реализация функции выше работает на Scala / Java API Sedona. При написании этой функции на Python API Sedona все получалось, кроме гибкого типа пространственного джойна. То, что в Python API его не было, меня не устраивало. Поэтому изначально я пытался сделать джарник с этой моей функцией и потом написать некую питоновскую обертку.
А потом понял, что все поведение кастомной функции со сложной логикой можно имитировать через настройки для Sedona в рамках конфига Spark-сессии. Ларчик просто открывался!
В итоге я подобрал оптимальные конфиги для коллег для решения данной задачи, которые были продуктивизированы.
Пройдемся по конфигам Spark-сессии:
# https://sedona.apache.org/1.7.2/api/sql/Parameter/ "sedona.global.index": "true" "sedona.global.indextype": "quadtree" "sedona.join.gridtype": "quadtree" "sedona.join.indexbuildside": "left" "sedona.join.numpartition": 12800 "sedona.join.spatitionside ": "right" "sedona.join.optimizationmode": "all" "sedona.join.autoBroadcastJoinThreshold": "-1"
sedona.global.index — включается индексация при джойнах.
sedona.global.indextype — указываем квадратное дерево (могли указать в кастомной функции).
sedona.join.gridtype — здесь также указываем квадратное дерево (могли указать в типе построения пространственной сетки).
sedona.join.indexbuildside — на левом или правом датафрейме при джоине строим пространственный индекс
sedona.join.numpartition — те же самые 12 800 ячеек. Напомню, это число ячеек при построении пространственной сетки, которую использует дальше Sedona для индексации джойна.
sedona.join.spatitionside. Это не опечатка, а сокращение от spatial partition side — сторона пространственного позиционирования в сокращенном виде. Здесь она указана как right, то есть применительно к нашей задаче с бенчмарками это правый датафрейм, в котором 11 млрд событий.
sedona.join.optimizationmode — включаем в Sedona режим оптимизации.
sedona.join.autoBroadcastJoinThreshold здесь -1.
Кажется, что Broadcast — это хорошо, почему же мы их выключили? Дело в том, что если посмотреть в Spark UI в раздел SQL/DataFrame, то увидим план запроса с нашими геометриями. У Spark есть свои типы физических джойнов: Sort Merge, Broadcast в реализациях Broadcast Hash и Broadcast Nested Loops, есть Shuffle Hash. А у Sedona еще и свои типы физических джойнов. В частности, в Sedona все эти оптимизации отлично работают на физическом типе Range Join, а на Sedona-версии Broadcast, к сожалению - нет. Поэтому мы выключаем Broadcast принудительно, используем Range Join для оптимизированного режима джойна.
С этими тонкими настройками на тех же ресурсах при статической аллокации мы получили очень обнадеживающие результаты.
Ресурсы — join Spark DataFrame с геоданными на Apache Sedona:
from sedona.utils import SedonaKryoRegistrator, KryoSerializer "spark.executor.cores": "4", "spark.executor.memory": "12g", "spark.dynamicAllocation.enabled": "false", "spark.executor.instances": "200", "spark.sql.shuffle.partitions": "800", "spark.default.parallelism": "800", "spark.jars.packages": [ "org.apache.sedona:sedona-spark-shaded-3.3_2.12:1.7.2", "org.datasyslab:geotools-wrapper:1.5.1-28.2", ], "spark.serializer": KryoSerializer.getName, "spark.kryo.registrator": SedonaKryoRegistrator.getName,

На тех же данных и на тех же ресурсах — join с тонкими настройками конфигурации SparkSession + дата-фрейм API дает те же результаты, что кастомный инструмент на Scala API Sedona + SpatialRDD со сложной логикой.
Код с джойном дата-фрейма:
from sedona.sql.st_predicates import ST_Contains # Benchmark # sdf_events - 11_770_557_003 rows # DataFrames with Data Skew > 1.5 hour # sdf_polygons - 72_650_028 rows # SpatialRDDs with Index = 15 min # DataFrames with Index = 15 min sdf_joined = ( sdf_polygons.alias("p") .join( sdf_events.alias("e"), ST_Contains(col("p.geometry"), col("e.point")), "inner" ) .select( col("e.hash"), col("e.time"), col("e.point").cast("string").alias("point"), col("p.region_name"), col("p.district_name") ) )
Мы видим, что по умолчанию с перекосом данных у нас не менее 1,5 часов происходит джойн, финальные Spark-таски висят — все плохо.
Без перекоса данных и с оптимизацией — что в случае кастомной функции, что в случае оптимальных настроек — мы получаем результат за 15 минут: как минимум выигрыш в шесть раз, и это на данных по всей России! Если перейти к данным на меньшей площади (например, я также исследовал по набору данных в Москве и Московской области), получаем еще больший выигрыш — более чем в 10 раз по времени.
Что в итоге
Выводы из этой истории я сделал следующие:
Часто разработка кастомных велосипедов совсем не бесполезна. Напротив, она позволяет добиться глубокого понимания процесса — хоть и набить себе шишки, но в итоге выйти на более элегантное и изящное решение, которое может быть продуктивизировано.
При джойнах на Apache Sedona для оптимальной производительности однозначно используйте или конфиг с тонкими настройками, или, если вы не ищете легких путей, что-то подобное кастомному инструменту.
Если вы занимаетесь обработкой геометрических данных и их джойнами, то нужно понимать, какие бывают типы пространственных сеток, режимы оптимизации, какие там могут быть типы индексированных джойнов.
При джойне стройте пространственную сетку по тому из дата-фреймов, данные которого распределены в пространстве неравномерно. Это поможет избежать Data Skew.
Проводите осмысленный подбор параметра numpartition, отвечающего за количество ячеек в пространственной сетке. Слишком низкое его значение снова приведет к Data Skew.