Predicate Pushdown в Object Storage: что пушится, что нет

Predicate Pushdown - водораздельная линия между эффективным кодом и кодом, который сжигает деньги в облаке. Разбираем механику, поддерживаемые и неподдерживаемые предикаты, диагностику через explain() и оптимизацию.

storage

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 (какие байты запросить).

Таким образом существуют два уровня фильтрации:

  1. Row Group Level (metadata filtering): предикат применяется к min/max статистике в Footer → пропуск целых Row Groups без чтения данных. Это настоящий Predicate Pushdown - данные вообще не читаются из S3.

  2. 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: Диагностировать причины

Два проблемных предиката:

  1. YEAR(created_at) = 2024 - функция над колонкой, PushedFilters = пустой.
  2. 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

  1. Никаких функций в WHERE над колонками, по которым нужна эффективная фильтрация. WHERE year(ts) = 2024WHERE ts >= '2024-01-01' AND ts < '2025-01-01'.

  2. Типы должны совпадать: col STRING → фильтр STRING. col BIGINT → фильтр INT64. Несовпадение типов = нет pushdown.

  3. OR - опасен: комбинируйте через UNION или добавляйте AND-условия, сужающие выборку.

  4. Python UDF в filter - почти всегда заменима встроенными функциями Spark. Если нельзя - хотя бы добавьте pushdown-фильтр рядом.

  5. LIKE '%text%' - не пушится. Для full-text search - Elasticsearch.

  6. Сортировка данных при записи: sortWithinPartitions(filter_col) делает Row Group Filtering эффективным.

  7. Z-ordering (Delta Lake) или Sort Order (Iceberg) для многомерной кластеризации.

  8. Проверяйте через explain() перед деплоем: PushedFilters должен содержать нужные условия, ReadSchema - минимальный набор колонок.

  9. Spark UI → SQL → FileScan: смотрите на bytes scanned и files pruned - они показывают реальную эффективность pushdown.

  10. Lakehouse-форматы (Delta/Iceberg) дают дополнительный уровень file-level statistics, позволяя пропускать файлы до открытия их Footer.