Predicate Pushdown в Object Storage: что пушится, что нет
Predicate Pushdown - водораздельная линия между эффективным кодом и кодом, который сжигает деньги в облаке. Разбираем механику, поддерживаемые и неподдерживаемые предикаты, диагностику через explain() и оптимизацию.
Scan Avoidance важнее быстрого чтения¶
В мире облачной аналитики существует контринтуитивное правило: не читать данные быстро - читать как можно меньше данных. Разница принципиальная.
Ускорить чтение из S3 почти невозможно - это ограничение физики: HTTP через публичную сеть, латентность 5–50 мс на запрос, пропускная способность в сотни MB/s. Но избежать чтения ненужных данных - можно и нужно. Именно это делает Predicate Pushdown.
Рассмотрим запрос к таблице событий размером 5 ТБ:
SELECT user_id, SUM(amount) as total
FROM events
WHERE event_date = '2024-06-15' AND country = 'RU'
GROUP BY user_id
Без Predicate Pushdown: Spark читает все 5 ТБ из S3, загружает в память экзекуторов, применяет фильтры в памяти. Итог: 5 ТБ сетевого трафика, 30 минут, $115 за egress.
С Predicate Pushdown: Spark читает Footer каждого Parquet-файла (метаданные), определяет, какие Row Groups содержат event_date = '2024-06-15', читает из S3 только их - ~50 ГБ. Итог: 100× меньше трафика, 2 минуты, $1.15.
Это не теоретическое преимущество - это разница между системой, которую команда готова поддерживать, и системой, которая поедает бюджет и уволила трёх аналитиков за «медленные запросы».
Архитектура Predicate Pushdown: как это работает¶
Путь фильтра от SQL до S3¶
Важно понимать: PushedFilters - это не значит, что фильтр выполняется «в S3». S3 - тупое хранилище объектов, оно не понимает SQL. Pushdown означает, что фильтр передаётся в Parquet Reader, который использует его для принятия решений на уровне метаданных (какие Row Groups пропустить) и HTTP Range Requests (какие байты запросить).
Таким образом существуют два уровня фильтрации:
-
Row Group Level (metadata filtering): предикат применяется к min/max статистике в Footer → пропуск целых Row Groups без чтения данных. Это настоящий Predicate Pushdown - данные вообще не читаются из S3.
-
Page/Row Level (in-memory filtering): после загрузки Row Group данные проходят точную фильтрацию уже в памяти экзекутора. Это тоже часть pushed filter - просто менее эффективная, так как данные уже прочитаны.
Partition Pruning vs Predicate Pushdown¶
Студенты часто путают два механизма. Это разные уровни оптимизации:
Partition Pruning - это уровень файловой системы. Если данные партиционированы по event_date, Spark при запросе WHERE event_date = '2024-06-15' вообще не заглядывает в другие партиции. Это происходит до открытия любого файла - используются только метаданные каталога.
Predicate Pushdown - уровень внутри файла. Даже если партиция event_date=2024-06-15 прочитана, файл может содержать 100 Row Groups, из которых нужны только 5. Predicate Pushdown пропускает остальные 95, проверив статистику в Footer.
Оба механизма должны работать вместе. Partition Pruning без Predicate Pushdown - как найти нужный шкаф, но обыскать все ящики. Predicate Pushdown без Partition Pruning - как искать по всей квартире, но умно пропускать нерелевантные ящики в каждой комнате.
Что успешно пушится: поддерживаемые предикаты¶
Операторы сравнения для примитивных типов¶
Наиболее эффективный pushdown - для простых условий сравнения на числовых типах, датах и строках:
from pyspark.sql import functions as F
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("pushdown-demo").getOrCreate()
events = spark.read.parquet("s3a://data-lake/events/")
# ✅ EqualTo - пушится идеально:
events.filter(F.col("status") == "PAID")
# PushedFilters: [IsNotNull(status), EqualTo(status,PAID)]
# ✅ LessThan / GreaterThan / LessThanOrEqual / GreaterThanOrEqual:
events.filter(F.col("amount") > 1000)
events.filter(F.col("amount").between(100, 500))
# PushedFilters: [IsNotNull(amount), GreaterThan(amount,1000.0)]
# PushedFilters: [GreaterThanOrEqual(amount,100.0), LessThanOrEqual(amount,500.0)]
# ✅ Даты и Timestamp:
events.filter(F.col("event_date") == "2024-06-15")
events.filter(F.col("created_at") >= "2024-06-01T00:00:00")
# PushedFilters: [IsNotNull(event_date), EqualTo(event_date,2024-06-15)]
# ✅ In - пушится как список EqualTo условий:
events.filter(F.col("country").isin("RU", "BY", "KZ"))
# PushedFilters: [In(country, [RU,BY,KZ])]
# ✅ IS NULL / IS NOT NULL - пушится через null_count статистику:
events.filter(F.col("phone").isNull())
events.filter(F.col("email").isNotNull())
# PushedFilters: [IsNull(phone)]
# PushedFilters: [IsNotNull(email)]
Почему примитивные типы работают хорошо: Parquet хранит в Footer точные числовые min и max для каждого Row Group. Проверить amount > 1000 - это одно сравнение с двумя числами из Footer. Никаких данных из S3 читать не нужно.
Логическое И (AND): комбинирование фильтров¶
AND-условия пушатся независимо и усиливают друг друга:
# ✅ AND: оба условия применяются к Row Group statsitics:
events.filter(
(F.col("event_date") == "2024-06-15") &
(F.col("amount") > 500) &
(F.col("country") == "RU")
)
# PushedFilters: [
# IsNotNull(event_date), EqualTo(event_date,2024-06-15),
# IsNotNull(amount), GreaterThan(amount,500.0),
# IsNotNull(country), EqualTo(country,RU)
# ]
# Row Group пропускается если ХОТЯ БЫ ОДНО условие ложно по статистике:
# - max(event_date) < 2024-06-15 → SKIP
# - OR: max(amount) <= 500 → SKIP
# - OR: country словарь не содержит RU → SKIP
# Все три фильтра работают параллельно, взаимно усиливая отсечение
Это мощная особенность: каждый дополнительный AND-предикат только увеличивает вероятность пропустить Row Group. Три условия через AND могут отсечь 99.9% Row Groups, даже если каждое по отдельности отсекало бы только 90%.
LIKE с префиксом: ограниченный pushdown¶
# ✅ LIKE с префиксом (без ведущего wildcard) - частичный pushdown:
events.filter(F.col("name").like("Иван%"))
# PushedFilters: [StringStartsWith(name,Иван)]
# Как работает: Parquet может проверить через словарь (dictionary encoding),
# есть ли в Row Group значения начинающиеся с "Иван"
# Эффективно для колонок с небольшим количеством уникальных значений
# ✅ StringContains - если не в начале, pushdown ограничен:
events.filter(F.col("name").contains("Иван"))
# PushedFilters: [StringContains(name,Иван)]
# Но эффективен ТОЛЬКО если column использует dictionary encoding
# ✅ Not - работает для простых условий:
events.filter(~(F.col("status") == "CANCELLED"))
# PushedFilters: [Not(EqualTo(status,CANCELLED))]
Boolean-колонки¶
# ✅ Boolean условия:
events.filter(F.col("is_premium") == True)
events.filter(F.col("is_deleted").cast("boolean"))
# PushedFilters: [IsNotNull(is_premium), EqualTo(is_premium,true)]
Проверка через explain()¶
Всегда проверяйте через .explain(), что нужные фильтры действительно попали в PushedFilters:
df = events \
.filter(F.col("event_date") == "2024-06-15") \
.filter(F.col("amount") > 500) \
.select("user_id", "amount")
df.explain()
# Ожидаемый вывод:
# == Physical Plan ==
# *(1) Project [user_id#0L, amount#1]
# +- *(1) Filter (isnotnull(amount#1) AND (amount#1 > 500.0))
# +- *(1) FileScan parquet [user_id#0L, amount#1, event_date#2]
# Batched: true
# Location: InMemoryFileIndex(...)
# PushedFilters: [IsNotNull(event_date), EqualTo(event_date,2024-06-15),
# IsNotNull(amount), GreaterThan(amount,500.0)]
# ReadSchema: struct<user_id:bigint, amount:double, event_date:date>
Обратите внимание: EqualTo(event_date,2024-06-15) есть в PushedFilters, значит Row Group Filtering будет применён. ReadSchema содержит только 3 колонки из N - Column Pruning работает.
Что НЕ пушится: ловушки, которые жгут деньги¶
Ловушка 1: Логическое ИЛИ (OR) на разных колонках¶
OR - самая опасная ловушка для Predicate Pushdown:
# ❌ OR на разных колонках - pushdown частично ломается:
events.filter(
(F.col("amount") > 10000) |
(F.col("is_fraud") == True)
)
# PushedFilters: [Or(GreaterThan(amount,10000.0), EqualTo(is_fraud,true))]
# Проблема: Row Group пропускается ТОЛЬКО если ОБА условия ложны по статистике:
# - max(amount) <= 10000 AND is_fraud словарь не содержит true → SKIP
# - Но если max(amount) > 10000 (а это почти всегда) → Row Group читается
# В реальности: OR почти всегда заставляет читать все Row Groups
Почему OR проблематичен: чтобы пропустить Row Group при OR, нужно доказать, что ОЕБА условия заведомо ложны для всех строк в группе. Это случается редко. Если хотя бы одно условие «может быть истинным» - Row Group читается.
# ✅ Переписать через UNION (если данные позволяют):
df1 = events.filter(F.col("amount") > 10000)
df2 = events.filter(F.col("is_fraud") == True)
result = df1.union(df2).distinct()
# Каждый запрос имеет свой pushdown,
# хотя в целом запросов к S3 будет больше
# ✅ Или через агрегацию IF:
# Вместо OR - использовать данные, которые удовлетворяют хотя бы одному
events.filter(
(F.col("amount") > 10000) | (F.col("is_fraud") == True)
).show()
# Если OR неизбежен - хотя бы добавьте дополнительные AND-условия,
# которые сузят выборку до читаемых Row Groups:
events.filter(
(F.col("event_date") == "2024-06-15") & # AND сужает через partition
((F.col("amount") > 10000) | (F.col("is_fraud") == True)) # OR внутри дня
)
Ловушка 2: Функции над колонками¶
# ❌ Любая функция над колонкой убивает pushdown для этой колонки:
events.filter(F.upper(F.col("status")) == "PAID")
# PushedFilters: [] ← пустой! Весь файл читается, потом фильтрация в памяти
events.filter(F.col("email").contains("@gmail.com"))
# PushedFilters: [StringContains(email,@gmail.com)]
# ← Частично пушится, но ONLY через dictionary (если dictionary encoding)
events.filter(F.year(F.col("created_at")) == 2024)
# PushedFilters: [] ← Функция year() не может быть применена к min/max Timestamp
events.filter(F.date_add(F.col("created_at"), 5) > F.current_date())
# PushedFilters: [] ← Вычисляемые выражения не пушатся
events.filter(F.col("price") * 1.2 > 1000)
# PushedFilters: [] ← Арифметика над колонкой отключает pushdown
Почему функции ломают pushdown: Parquet Footer содержит min/max для «сырых» значений колонки. Для status в Footer может быть min='active', max='paid'. Но что такое upper(min('active'))? Parquet Reader не знает - ему пришлось бы выполнить функцию над каждым значением, а это значит - читать все данные. Проще не пытаться.
# ✅ Переписать функции в статические условия:
# ❌ Было:
events.filter(F.upper(F.col("status")) == "PAID")
# ✅ Стало: все варианты написания явно
events.filter(F.col("status").isin("paid", "PAID", "Paid"))
# PushedFilters: [In(status, [paid,PAID,Paid])]
# ❌ Было:
events.filter(F.year(F.col("created_at")) == 2024)
# ✅ Стало: диапазон для Timestamp:
events.filter(
(F.col("created_at") >= "2024-01-01T00:00:00") &
(F.col("created_at") < "2025-01-01T00:00:00")
)
# PushedFilters: [GreaterThanOrEqual(created_at,...), LessThan(created_at,...)]
# ❌ Было:
events.filter(F.col("price") * 1.2 > 1000)
# ✅ Стало: перенести математику на константу:
events.filter(F.col("price") > 1000 / 1.2) # = 833.33
# PushedFilters: [GreaterThan(price,833.3333...)]
Ловушка 3: LIKE с ведущим wildcard¶
# ❌ LIKE с подстановкой в начале - pushdown невозможен:
events.filter(F.col("email").like("%gmail.com"))
# PushedFilters: [] ← Parquet не может проверить суффикс через min/max
events.filter(F.col("description").like("%ошибка%"))
# PushedFilters: [] ← Поиск подстроки невозможен через метаданные Parquet
# ✅ Альтернативы:
# 1. Если возможно - хранить домен отдельной колонкой:
events.filter(F.col("email_domain") == "gmail.com")
# PushedFilters: [EqualTo(email_domain,gmail.com)]
# 2. Для full-text поиска - использовать Elasticsearch/OpenSearch,
# а не Spark-запросы к Parquet на S3
Ловушка 4: Несоответствие типов (Type Mismatch)¶
# ❌ Колонка STRING, фильтр NUMBER - pushdown отключается:
# Схема: phone_number STRING
events.filter(F.col("phone_number") == 79991112233) # число!
# PushedFilters: [] ← Нужно cast, cast убивает pushdown
# ✅ Фильтровать с правильным типом:
events.filter(F.col("phone_number") == "79991112233") # строка!
# PushedFilters: [EqualTo(phone_number,79991112233)]
# ❌ Явный cast - тоже убивает pushdown:
events.filter(F.col("user_id").cast("string") == "12345")
# PushedFilters: []
# ✅ Сравниваем с типом колонки:
events.filter(F.col("user_id") == 12345) # user_id BIGINT
# PushedFilters: [EqualTo(user_id,12345)]
# ❌ Timestamp vs String - тонкий type mismatch:
events.filter(F.col("created_at") == "2024-06-15") # String vs Timestamp
# PushedFilters: [] или некорректный pushdown
# ✅ Использовать to_timestamp / to_date:
events.filter(F.col("event_date") == F.lit("2024-06-15").cast("date"))
# или через SQL:
spark.sql("SELECT * FROM events WHERE event_date = DATE '2024-06-15'")
# PushedFilters: [EqualTo(event_date,2024-06-15)]
Ловушка 5: Фильтры после трансформации¶
# ❌ Фильтр применяется ПОСЛЕ join или withColumn - данные уже в памяти:
result = events \
.join(users, "user_id") \
.withColumn("adjusted_amount", F.col("amount") * 1.1) \
.filter(F.col("adjusted_amount") > 1000)
# PushedFilters для events: [] (filter на вычисляемом поле)
# Spark читает все данные events, потом фильтрует в памяти
# ✅ Применить фильтр ДО join (с оригинальным полем):
events_filtered = events.filter(F.col("amount") > 1000 / 1.1) # ≈ 909
result = events_filtered \
.join(users, "user_id") \
.withColumn("adjusted_amount", F.col("amount") * 1.1)
# PushedFilters для events: [GreaterThan(amount,909.09...)]
# Spark читает только события с amount > 909 → потом join
Ловушка 6: Python UDF в filter¶
from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType
# ❌ Python UDF в filter - полный отказ от pushdown:
@udf(returnType=BooleanType())
def is_valid_phone(phone: str) -> bool:
return phone is not None and len(phone) == 11 and phone.startswith("7")
events.filter(is_valid_phone(F.col("phone")))
# PushedFilters: [] ← UDF непрозрачен для Catalyst Optimizer
# Spark читает ALL данные, передаёт в Python, выполняет UDF построчно
# → 10–100× медленнее чем встроенные функции
# ✅ Заменить UDF на встроенные функции Spark:
events.filter(
F.col("phone").isNotNull() &
(F.length(F.col("phone")) == 11) &
F.col("phone").startswith("7")
)
# PushedFilters: [IsNotNull(phone), StringStartsWith(phone,7)]
# Часть pushdown'ится, остальное выполняется в векторизованном Java-коде
# (намного быстрее Python UDF)
Ловушка 7: Фильтрация внутри сложных структур¶
# ❌ Фильтр по элементу ARRAY - не пушится:
events.filter(F.array_contains(F.col("tags"), "promo"))
# PushedFilters: [] ← Parquet не хранит статистику по элементам массива
# ❌ Фильтр по ключу MAP - не пушится:
events.filter(F.col("metadata")["source"] == "mobile")
# PushedFilters: []
# ❌ Фильтр по глубоко вложенному STRUCT - ограниченный pushdown:
events.filter(F.col("user.address.city") == "Moscow")
# PushedFilters: [EqualTo(user.address.city,Moscow)] - может пушиться
# (зависит от версии Spark и глубины вложенности)
# ✅ Денормализация для часто используемых фильтров:
# Вместо tags ARRAY<STRING> - хранить отдельные boolean колонки:
# is_promo BOOLEAN, is_vip BOOLEAN, is_mobile BOOLEAN
events.filter(F.col("is_promo") == True)
# PushedFilters: [EqualTo(is_promo,true)] - пушится отлично
Визуальная карта поддержки предикатов¶
Влияние физической организации данных на Pushdown¶
Почему неупорядоченные данные делают Pushdown бесполезным¶
Predicate Pushdown проверяет min/max статистику Row Groups. Но если данные не отсортированы - min/max практически одинаковы во всех Row Groups:
Неотсортированные данные (случайный порядок):
RG1: amount min=0.01, max=99999.99 (все суммы вперемешку)
RG2: amount min=0.50, max=98765.43
RG3: amount min=1.00, max=97654.32
Запрос: WHERE amount > 50000
→ Проверяем: max(RG1)=99999 > 50000 → READ
→ Проверяем: max(RG2)=98765 > 50000 → READ
→ Проверяем: max(RG3)=97654 > 50000 → READ
→ Все Row Groups читаются! Pushdown не помог ни разу.
Отсортированные данные (ORDER BY amount):
RG1: amount min=0.01, max=15234.56 (маленькие суммы)
RG2: amount min=15234.57, max=43210 (средние суммы)
RG3: amount min=43211, max=99999.99 (большие суммы)
Запрос: WHERE amount > 50000
→ Проверяем: max(RG1)=15234 < 50000 → SKIP
→ Проверяем: max(RG2)=43210 < 50000 → SKIP
→ Проверяем: max(RG3)=43211 ≤ max=99999 → READ
→ Читаем только 1 из 3 Row Groups! Pushdown эффективен.
Сортировка при записи: основной инструмент¶
from pyspark.sql import functions as F
def write_optimized_for_pushdown(
df,
output_path: str,
sort_cols: list,
partition_col: str = None
) -> None:
"""
Записать DataFrame с оптимальным физическим порядком
для эффективного Predicate Pushdown.
sort_cols: список колонок в порядке приоритета фильтрации.
Первая колонка - самая часто используемая в WHERE.
"""
writer = df \
.sortWithinPartitions(*sort_cols) \
.write \
.option("parquet.block.size", 268435456) # 256 MB Row Groups
if partition_col:
writer = writer.partitionBy(partition_col)
writer.mode("overwrite").parquet(output_path)
# Пример: таблица событий, часто фильтруемая по event_date и user_id:
events = spark.read.parquet("s3a://data-lake/events-raw/")
write_optimized_for_pushdown(
df=events,
output_path="s3a://data-lake/events-optimized/",
sort_cols=["event_date", "user_id"],
partition_col="event_date"
)
# Результат:
# - Partition Pruning по event_date (уровень директорий)
# - Внутри партиции: данные отсортированы по user_id
# → Row Group Filtering по user_id работает эффективно
# - WHERE event_date='2024-06-15' AND user_id BETWEEN 1000 AND 2000
# → пропускает большинство Row Groups
Z-ordering в Delta Lake: многомерная кластеризация¶
Сортировка работает только для одной колонки хорошо. Для нескольких колонок одновременно нужна Z-ordering:
# Z-Order оптимизирует сразу по нескольким измерениям:
spark.sql("""
OPTIMIZE delta.`s3a://data-lake/events`
ZORDER BY (user_id, event_date)
""")
# После Z-ordering:
# Запрос: WHERE user_id = 12345 AND event_date = '2024-06-15'
# → Row Groups содержат близкие значения ОБОИХ полей одновременно
# → Большинство RG пропускается по min/max обоих полей
# Проверить эффект в explain:
spark.sql("""
SELECT user_id, SUM(amount)
FROM delta.`s3a://data-lake/events`
WHERE user_id = 12345 AND event_date = '2024-06-15'
GROUP BY user_id
""").explain()
# PushedFilters: [EqualTo(user_id,12345), EqualTo(event_date,2024-06-15)]
Hidden Partitioning в Iceberg: автоматическая кластеризация¶
# Iceberg Hidden Partitioning автоматически применяет partition transforms:
spark.sql("""
CREATE TABLE catalog.events (
user_id BIGINT,
amount DOUBLE,
event_type STRING,
created_at TIMESTAMP
)
USING iceberg
PARTITIONED BY (
days(created_at), -- partition по дню
bucket(16, user_id) -- bucket по user_id для равномерного распределения
)
""")
# Запрос:
spark.sql("""
SELECT user_id, SUM(amount)
FROM catalog.events
WHERE created_at >= '2024-06-15'
AND created_at < '2024-06-16'
AND user_id = 12345
GROUP BY user_id
""")
# Iceberg выполняет:
# 1. Partition pruning: читает только snapshot с days=19889 (2024-06-15)
# 2. Bucket pruning: hash(12345) % 16 = bucket 9 → читаем только bucket 9
# 3. Row Group Filtering внутри файлов через Parquet Footer
# → 3 уровня отсечения данных вместо 2
Predicate Pushdown в Lakehouse-форматах¶
Дополнительный уровень: file-level statistics¶
Delta Lake и Iceberg хранят min/max не только в Parquet Footer каждого файла, но и в собственных метаданных - Delta Log / Manifest Files. Это позволяет пропускать целые файлы до их открытия:
Настройка статистик в Delta Lake¶
# Delta Lake собирает статистику при записи (включено по умолчанию):
df.write \
.format("delta") \
.option("delta.dataSkippingNumIndexedCols", "32") \ # статистика для первых 32 колонок
.mode("append") \
.save("s3a://data-lake/delta/events/")
# Проверить собранную статистику:
spark.sql("""
SELECT
path,
stats.numRecords,
stats.minValues.amount as min_amount,
stats.maxValues.amount as max_amount
FROM delta.`s3a://data-lake/delta/events`.files
""")
# Ручная пересборка статистики (если данные были записаны без неё):
spark.sql("""
ANALYZE TABLE delta.`s3a://data-lake/delta/events`
COMPUTE STATISTICS FOR ALL COLUMNS
""")
Настройка статистик в Iceberg¶
# Iceberg собирает statistics при записи через Write Metrics Mode:
spark.sql("""
ALTER TABLE catalog.events SET TBLPROPERTIES (
'write.metadata.metrics.mode' = 'full', -- full / counts / none
'write.metadata.metrics.column.user_id' = 'full', -- per-column override
'write.metadata.metrics.column.amount' = 'full',
'write.metadata.metrics.column.tags' = 'counts' -- только null counts для ARRAY
)
""")
# Проверить статистику в manifest files:
spark.sql("""
SELECT file_path, record_count, column_sizes, lower_bounds, upper_bounds
FROM catalog.events.files
ORDER BY record_count DESC
LIMIT 10
""")
Диагностика: находим проблемы в реальных запросах¶
Инструментарий диагностики¶
def analyze_pushdown(df, name: str = "Query") -> None:
"""
Диагностика Predicate Pushdown для DataFrame.
Выводит PushedFilters, ReadSchema и признаки неэффективности.
"""
import re
plan = df._jdf.queryExecution().simpleString()
print(f"\n{'='*60}")
print(f"Анализ: {name}")
print(f"{'='*60}")
# Найти PushedFilters:
pushed = re.findall(r"PushedFilters: \[(.*?)\]", plan)
if pushed:
filters = pushed[0]
print(f"✅ PushedFilters: [{filters}]")
if not filters:
print(" ⚠️ PushedFilters пустой - фильтры не пушатся!")
else:
print("❌ PushedFilters не найден - проверьте формат файла")
# Найти ReadSchema:
schema = re.findall(r"ReadSchema: (struct<.*?>)", plan)
if schema:
cols = schema[0].count(",") + 1
print(f"📋 ReadSchema: {cols} колонок читается")
if "struct<" in schema[0] and cols > 20:
print(" ⚠️ Много колонок - проверьте Column Pruning")
# Проверить Batched:
if "Batched: true" in plan:
print("✅ Vectorized Reader: активен (Batched: true)")
elif "Batched: false" in plan:
print("⚠️ Vectorized Reader: отключён (Batched: false)")
print()
# Использование:
df_bad = events.filter(F.upper(F.col("status")) == "PAID")
analyze_pushdown(df_bad, "Запрос с upper() функцией")
df_good = events.filter(F.col("status").isin("paid", "PAID"))
analyze_pushdown(df_good, "Оптимизированный запрос")
Мониторинг через Spark UI¶
В Spark UI (вкладка SQL → запрос → план выполнения) найдите оператор FileScan и посмотрите метрики:
FileScan parquet (метрики оператора):
├── number of files read: 27 - файлов прочитано
├── files pruned: 9,973 - файлов пропущено partition pruning
├── scan time total: 00:01:23 - время чтения S3
├── metadata time: 00:00:05 - время чтения Footer/метаданных
├── number of output rows: 45,000 - строк прошло через scan
└── bytes scanned: 850,341,888 - байт прочитано из S3 (~850 MB)
Диагностика:
- bytes scanned / number of output rows = 850MB / 45K = ~19KB/row
Если avg row size = 1KB → читаем в 19× больше нужного → плохой pushdown
- files pruned = 9,973 из 10,000 → хорошо! 99.7% файлов пропущено
- scan time >> metadata time → хорошо (не много мелких запросов к Footer)
- scan time << job time → плохо если scan time = 80% job time
Anti-patterns: диагностика и исправление¶
Антипаттерн 1: «Pushdown работает, но данные перемешаны»¶
# Симптом: PushedFilters заполнен, но bytes_scanned огромный
# Диагностика:
events.filter(F.col("user_id") == 12345).explain()
# PushedFilters: [IsNotNull(user_id), EqualTo(user_id,12345)] ← есть!
# Но Spark UI показывает: bytes_scanned = 950 GB из 1 TB
# Причина: данные не кластеризованы по user_id
# Min/max каждой Row Group: 1 ≤ user_id ≤ 99999999 (почти все пользователи)
# → ни одна Row Group не пропускается, хотя фильтр «пушится»
# Исправление: добавить сортировку при записи:
events.repartition(200, F.col("user_id")) \
.sortWithinPartitions("user_id") \
.write \
.option("parquet.block.size", 268435456) \
.mode("overwrite") \
.parquet("s3a://data-lake/events-by-userid/")
# Теперь Row Groups содержат узкие диапазоны user_id
# WHERE user_id = 12345 → читает только несколько Row Groups
Антипаттерн 2: Фильтр через join condition вместо where¶
# ❌ Join condition не даёт pushdown на большую таблицу:
events \
.join(
target_dates.filter(F.col("date") == "2024-06-15"),
events.event_date == target_dates.date
)
# Spark может не «видеть», что events нужно фильтровать по event_date
# ✅ Явный фильтр перед join:
events.filter(F.col("event_date") == "2024-06-15") \
.join(
target_dates.filter(F.col("date") == "2024-06-15"),
events.event_date == target_dates.date
)
# PushedFilters для events: [EqualTo(event_date,2024-06-15)]
Антипаттерн 3: filter после cache()¶
# ❌ filter после cache - данные уже в памяти, pushdown был до cache:
events_cached = events.cache()
events_cached.filter(F.col("country") == "RU") # filter после cache
# Pushdown НЕ применяется к данным в кэше
# Зато они уже в памяти - фильтрация быстрая (в памяти, не S3)
# НО: при cache() Spark читает ВСЕ данные из S3 → overhead при кешировании
# ✅ Применять фильтры ДО cache (только нужные данные попадают в кэш):
events_ru = events \
.filter(F.col("country") == "RU") \
.cache()
# PushedFilters при populate cache: [EqualTo(country,RU)]
# → В кэш попадает только ~10% данных
Антипаттерн 4: Игнорирование stale statistics¶
# После частых appends в Delta/Iceberg статистика может устареть:
# Новые файлы не имеют column statistics → Spark читает их полностью
# ✅ Регулярно обновлять статистику:
spark.sql("ANALYZE TABLE catalog.events COMPUTE STATISTICS FOR ALL COLUMNS")
# ✅ Для Delta Lake: OPTIMIZE пересобирает файлы со свежей статистикой:
spark.sql("OPTIMIZE delta.`s3a://data-lake/events`")
# ✅ Проверить, есть ли статистика в файлах:
spark.sql("""
SELECT
file_path,
CASE WHEN lower_bounds IS NULL THEN 'NO STATS' ELSE 'HAS STATS' END as stats_status
FROM catalog.events.files
LIMIT 20
""")
Production-кейс: отладка медленного Spark-запроса¶
Разберём реальный сценарий: запрос к таблице заказов выполняется 40 минут вместо ожидаемых 2 минут.
Шаг 1: Собрать симптомы¶
# Исходный «медленный» запрос:
result = spark.sql("""
SELECT
DATE_FORMAT(created_at, 'yyyy-MM') as month,
status,
COUNT(*) as cnt,
SUM(amount) as total
FROM orders
WHERE YEAR(created_at) = 2024
AND UPPER(status) IN ('COMPLETED', 'PAID')
GROUP BY 1, 2
""")
result.explain()
# == Physical Plan ==
# *(2) HashAggregate(keys=[date_format(...)], functions=[count, sum])
# +- Exchange hashpartitioning(...)
# +- *(1) HashAggregate(keys=[date_format(...)], functions=[partial_count, partial_sum])
# +- *(1) Project [...]
# +- *(1) Filter ((year(created_at) = 2024) AND
# upper(status) IN ('COMPLETED', 'PAID'))
# +- *(1) FileScan parquet [created_at, status, amount]
# Batched: true
# PushedFilters: [] ← ПУСТО! Это проблема.
# ReadSchema: struct<created_at:timestamp, status:string, amount:double>
# Spark UI показывает:
# bytes scanned: 4.7 TB ← читается почти вся таблица
# scan time: 38 minutes
# output rows after scan: 12,450,000
Шаг 2: Диагностировать причины¶
Два проблемных предиката:
YEAR(created_at) = 2024- функция над колонкой,PushedFilters= пустой.UPPER(status) IN (...)- функцияUPPERубивает pushdown для status.
Шаг 3: Переписать запрос¶
# ✅ Оптимизированный запрос:
result_opt = spark.sql("""
SELECT
DATE_FORMAT(created_at, 'yyyy-MM') as month,
status,
COUNT(*) as cnt,
SUM(amount) as total
FROM orders
WHERE created_at >= TIMESTAMP '2024-01-01 00:00:00'
AND created_at < TIMESTAMP '2025-01-01 00:00:00'
AND status IN ('completed', 'COMPLETED', 'paid', 'PAID')
GROUP BY 1, 2
""")
result_opt.explain()
# PushedFilters: [
# GreaterThanOrEqual(created_at,2024-01-01 00:00:00),
# LessThan(created_at,2025-01-01 00:00:00),
# In(status, [completed,COMPLETED,paid,PAID])
# ]
# bytes scanned: 420 GB (было: 4.7 TB → 11× меньше)
# scan time: 3 minutes 24 seconds (было: 38 минут)
Шаг 4: Улучшить физическую организацию данных¶
# Если запросы часто по created_at - переписать с партиционированием:
spark.sql("""
CREATE TABLE orders_v2
USING parquet
PARTITIONED BY (order_month)
AS
SELECT
*,
DATE_FORMAT(created_at, 'yyyy-MM') as order_month
FROM orders
SORT BY (created_at, status)
""")
# Или через DataFrame API с правильным порядком:
orders = spark.read.parquet("s3a://data-lake/orders/")
orders \
.withColumn("order_month", F.date_format("created_at", "yyyy-MM")) \
.repartition(F.col("order_month")) \
.sortWithinPartitions("created_at", "status") \
.write \
.partitionBy("order_month") \
.option("parquet.block.size", 268435456) \
.mode("overwrite") \
.parquet("s3a://data-lake/orders-optimized/")
# Запрос к оптимизированной таблице:
spark.read.parquet("s3a://data-lake/orders-optimized/") \
.filter(
(F.col("order_month") >= "2024-01") &
(F.col("order_month") <= "2024-12") &
F.col("status").isin("completed", "COMPLETED", "paid", "PAID")
) \
.groupBy(F.col("order_month"), F.col("status")) \
.agg(F.count("*").alias("cnt"), F.sum("amount").alias("total")) \
.explain()
# PushedFilters: [
# GreaterThanOrEqual(order_month,2024-01),
# LessThanOrEqual(order_month,2024-12),
# In(status, [completed,COMPLETED,paid,PAID])
# ]
# + Partition Pruning: только 12 партиций из возможных 60
# Результат:
# bytes scanned: 35 GB (из 4.7 TB → 134× меньше чем исходный запрос)
# scan time: 18 секунд
Финальное сравнение¶
| Версия | bytes scanned | Scan time | PushedFilters |
|---|---|---|---|
| Исходный (YEAR + UPPER) | 4.7 TB | 38 мин | [] - пустой |
| Переписанный предикат | 420 GB | 3.5 мин | 3 фильтра |
| + Партиционирование + сортировка | 35 GB | 18 сек | 3 фильтра + Partition Pruning |
Итого: в 134 раза меньше трафика, в 127 раз быстрее. Это не оптимизация «под капотом» - это разница в понимании того, как Predicate Pushdown взаимодействует с физическим форматом данных.
Чеклист: правила Predicate Pushdown для production¶
-
Никаких функций в WHERE над колонками, по которым нужна эффективная фильтрация.
WHERE year(ts) = 2024→WHERE ts >= '2024-01-01' AND ts < '2025-01-01'. -
Типы должны совпадать:
col STRING→ фильтрSTRING.col BIGINT→ фильтрINT64. Несовпадение типов = нет pushdown. -
OR - опасен: комбинируйте через UNION или добавляйте AND-условия, сужающие выборку.
-
Python UDF в filter - почти всегда заменима встроенными функциями Spark. Если нельзя - хотя бы добавьте pushdown-фильтр рядом.
-
LIKE '%text%' - не пушится. Для full-text search - Elasticsearch.
-
Сортировка данных при записи:
sortWithinPartitions(filter_col)делает Row Group Filtering эффективным. -
Z-ordering (Delta Lake) или Sort Order (Iceberg) для многомерной кластеризации.
-
Проверяйте через explain() перед деплоем:
PushedFiltersдолжен содержать нужные условия,ReadSchema- минимальный набор колонок. -
Spark UI → SQL → FileScan: смотрите на
bytes scannedиfiles pruned- они показывают реальную эффективность pushdown. -
Lakehouse-форматы (Delta/Iceberg) дают дополнительный уровень file-level statistics, позволяя пропускать файлы до открытия их Footer.