TGViewer
Поступашки — BIG DATA Поступашки — BIG DATA @bigdata_postupashki · 1.03K subscribers
Post #23 3.6K
Как Spark оптимизирует запросы

Spark вроде бы понятный — прочитал parquet, отфильтровал по статусу, посчитал count — всё просто. Но иногда одну и ту же команду можно написать так, что отработает за секунду, а в другом случае будет работать часы.
Почему? Давайте разберемся
Spark не магическая коробка, которая непонятно как молотит сотни гигабайт данных, а использует конкретные хорошо известные нам техники.
Я бы хотел остановиться на трех из них и показать, как не сломать их

1. Predicate pushdown
Spark может передавать фильтры (WHERE) сразу в источник — parquet, ORC, JDBC. То есть не грузить весь файл, а сразу читать только нужные строки.
Пример:
df = spark.read.parquet("s3://...").filter("status = 'active'")

Когда работает:
— фильтр простой (=, >, <)
— источник поддерживает pushdown (паркет, ORC, JDBC)
— фильтр указан явно, а не через UDF

Когда не работает:
—фильтр через UDF
—фильтр через выражение (df["a"] + df["b"] > 5)
—JSON и CSV (там pushdown нет вообще)
Проверить можно через план запроса df.explain(True). Ищите PushedFilters.

2. Column pruning
Если тебе нужны только user_id и amount, Spark может прочитать только эти поля. Особенно важно, если таблица на 300+ колонок.
Пример:
df = df.select("user_id", "amount")

Когда работает:
— n.select(...), без *
— формат — parquet, ORC, JDBC
Когда не работает:
— select("*")
— ты передал df в .rdd.map(...), .apply(...) или UDF
Поэтому, когда работаете с большими данными, рекомендую первое время вообще не использовать SELECT *, потому что в таблице может быть 100/200/100500 колонок, и тянуть их все влечет проблемы с производительностью. Вряд ли вам нужны все сто колонок
В explain будет Scan parquet [user_id, amount], если всё хорошо.

3. Partition pruning
Последнее по списку, но не по значимости
Если данные разложены по папкам (event_date=...), Spark может читать только нужные партиции, а не все подряд.
Пример:
df = spark.read.parquet("s3://events/").filter("event_date = '2025-07-30'")

аналогично через Spark SQL
spark.sql("SELECT id, date FROM table WHERE date = '20250801'")


Когда работает:
Очевидно ,когда таблица партицированна
Когда не работает:
— Фильтр через UDF или через функцию, например: year(event_time)
— Spark (да и не только) не умеет применять оптимизации к колонкам внутри функций
— JSON/CSV без схемы и partition discovery
Проверяется через PartitionFilters в explain(True).

Итак, точно ломают Spark-оптимизации:
— select("*"). Всегда.
— отсутствие фильтра по колонке партицирования
— функции на колонках, по которым идет фильтр
— UDF. В целом их можно использовать, но очень аккуратно и желательно стараться обходиться стандартными средствами
— использование формата вроде csv, json и т.д.

Так что следуйте этим простым правилам, и Spark будет у вас летать 😎
Ну и конечно всегда открывайте и читайте план запроса!

Подписывайтесь на канал, ставьте огонечки, будем и дальше разбираться с технологиями в Big Data! ⚡️
Также записывайтесь на наш курс по инженерии данных, чтобы стать крутым DE 🔥

@bigdata_postupashki
  • 🔥 9
  • 👍 2
More from @bigdata_postupashki
  1. Sep 6, 2026Разбор экзамена на стажировку в Т-банк за подписку! Чтобы получить разбор: ➡️Подпишитесь н…
  2. Jun 3, 2026Финальные скидки на карьерные курсы! Товарищи, специально для вас открываем доступ на лине…
  3. May 10, 2026Топ мест для вката в Data Engineering На слуху Яндекс, Т-банк, Авито и прочие ребята, но е…
  4. Apr 13, 2026Ментор+ Эта услуга для тех, кто хочет трудоустроиться уже сейчас без лишней головной боли…
  5. Jan 29, 2026Разбор экзамена по SQL Т-банка Разбор всех направлений выложен только на наших курсах. До…
  6. Dec 3, 2025Товарищи, Поступашкам нужны контент мейкеры. Если вы творческая личность, интересующейся ф…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →