Кэш без стратегии: 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() запланирован после последнего использования
✓ Кэш не создаётся в цикле без причины