Кэш без стратегии: persist до фильтрации, нет unpersist и cache в цикле

Кэширование - дорогая операция. Три классических ошибки: кэшировать сырые данные, не освобождать память и кэшировать в каждой итерации цикла.

optimization

Анти-паттерн 1: cache() до фильтрации

# ❌ Кэшируем весь сырой датасет - 1 ТБ в памяти кластера
raw_df = spark.read.parquet("s3://events/").cache()
raw_df.count()  # прогрев кэша

# Потом используем только 1% данных
result = raw_df.filter("event_type = 'purchase' AND country = 'RU'")
result.write.parquet("output/")

# Кэш занимает 1 ТБ RAM, используется для 10 ГБ данных
# ✅ Кэшируем ПОСЛЕ фильтрации
filtered_df = spark.read.parquet("s3://events/") \
    .filter("event_type = 'purchase' AND country = 'RU'") \
    .cache()
filtered_df.count()  # прогрев

# Теперь кэш = ~10 ГБ вместо 1 ТБ
result1 = filtered_df.groupBy("user_id").agg(sum("amount"))
result2 = filtered_df.join(users_df, "user_id")
filtered_df.unpersist()

Анти-паттерн 2: нет unpersist() → утечка памяти

# ❌ Кэш создаётся, но никогда не освобождается
def process_partition(date):
    df = spark.read.parquet(f"s3://events/date={date}").cache()
    df.count()
    return df.groupBy("user_id").sum("amount")

for date in all_dates:    # 365 итераций
    result = process_partition(date)
    result.write.parquet(f"s3://output/date={date}")
    # df.unpersist() ЗАБЫТО - кэш от каждой даты остаётся в памяти
    # После 10-20 итераций - OOM или вытеснение старых партиций + замедление
# ✅ Явный unpersist после использования
def process_partition(date):
    df = spark.read.parquet(f"s3://events/date={date}").cache()
    df.count()
    result = df.groupBy("user_id").sum("amount")
    df.unpersist()   # ← обязательно
    return result

Анти-паттерн 3: cache() в цикле итераций (ML anti-pattern)

# ❌ Кэш пересоздаётся в каждой итерации ML
for epoch in range(100):
    train_df = features_df \
        .withColumn("noise", rand())   # новый столбец каждый раз
        .cache()
    train_df.count()
    model.fit(train_df)
    train_df.unpersist()
# Каждый epoch - новый физический план, новый кэш → 100 cache/unpersist операций
# ✅ Кэшировать один раз ДО цикла
features_cached = features_df.cache()
features_cached.count()   # прогрев

for epoch in range(100):
    # Добавляем шум внутри model.fit(), не в кэшированном DF
    model.fit(features_cached)

features_cached.unpersist()   # после всех итераций

Анти-паттерн 4: cache() для одноразового датасета

# ❌ Кэш DF, который используется только один раз
df = spark.read.parquet("s3://...").cache()
df.write.parquet("output/")   # единственный Action → кэш бесполезен
# Время на запись кэша потрачено впустую

# ✅ cache() только если 2+ Actions на одном DF
if uses_count >= 2:
    df = df.cache()

Анти-паттерн 5: persist(MEMORY_ONLY) для больших данных

# ❌ MEMORY_ONLY при нехватке RAM - партиции молча отбрасываются и пересчитываются
df.persist(StorageLevel.MEMORY_ONLY)  # если не влезло - пересчёт, не ошибка

# ✅ MEMORY_AND_DISK - безопаснее: переполнение → диск, не потеря
df.persist(StorageLevel.MEMORY_AND_DISK)  # = df.cache()

# ✅ При тайтинге RAM - сериализованный вариант
df.persist(StorageLevel.MEMORY_AND_DISK_SER)  # Kryo compression → меньше RAM

Диагностика в Spark UI → Storage

RDD / DataFrame     Size in Memory   Size on Disk   Cached Partitions
events_filtered     8.2 GB           0 B            200/200    ← OK
raw_events          982 GB           12 GB          180/200    ← ПОДОЗРИТЕЛЬНО
old_daily_cache     45 GB            0 B            100/100    ← утечка, если не нужен
# Посмотреть текущий storage level программно
print(df.storageLevel)   # StorageLevel(True, True, False, True, 1)

# Очистить весь кэш (только для отладки!)
spark.catalog.clearCache()

Чек-лист кэширования

Before cache():
  ✓ DF используется 2+ раз?
  ✓ Данные отфильтрованы/агрегированы перед кэшем?
  ✓ Объём помещается в память кластера?

After cache():
  ✓ Явный count() или другой action для прогрева
  ✓ unpersist() запланирован после последнего использования
  ✓ Кэш не создаётся в цикле без причины