Speculative Execution: straggler detection и когда отключать
Spark запускает дублирующую Task для медленных стрэглеров. Разбираем механизм обнаружения, конфигурацию, идемпотентность и случаи, когда speculation делает хуже.
Проблема хвостовой латентности в распределённых системах¶
В Stage из 200 Tasks 199 завершаются за 30 секунд. Одна Task работает 15 минут. Весь Stage - и весь Job - ждёт эту одну. Это явление называется tail latency или straggler (отстающий).
В распределённых системах это системная проблема: чем больше узлов и Task'ей, тем выше вероятность, что хотя бы один окажется медленным. При 1000 узлах вероятность хотя бы одного сбоя в час приближается к 100%.
Причины страглеров: инфраструктурные vs. данные¶
Важно различать два класса причин - от этого зависит, поможет ли Speculative Execution.
Ключевое различие:
- Инфраструктурные причины - случайные. Перезапуск Task'и на другом Executor'е уберёт проблему. Speculative Execution помогает.
- Причины, связанные с данными - детерминированные. Дублирующая Task получит те же данные и будет такой же медленной. Speculative Execution не помогает.
Механика Speculative Execution¶
Алгоритм обнаружения страглеров¶
Параметры конфигурации¶
# Главный рубильник (по умолчанию false в OSS Spark, true в Databricks)
spark.conf.set("spark.speculation", "true")
# Как часто проверять наличие страглеров
spark.conf.set("spark.speculation.interval", "100ms") # default: 100ms
# Во сколько раз медленнее медианы должна быть Task
spark.conf.set("spark.speculation.multiplier", "1.5") # default: 1.5
# Какой % Tasks должен завершиться перед анализом
spark.conf.set("spark.speculation.quantile", "0.75") # default: 0.75
# Минимальное время работы Task перед проверкой
spark.conf.set("spark.speculation.minTaskRuntime", "100ms") # default: 100ms
Математика запуска: пример¶
Stage: 200 Tasks
Медиана завершённых Tasks: 30 сек
multiplier: 1.5
quantile: 0.75
Условия для запуска speculative copy Task 42:
1. Завершено >= 200 × 0.75 = 150 Tasks ✓
2. Task 42 работает > 30 × 1.5 = 45 сек ✓
3. Task 42 работает >= 100ms ✓
→ Spark запускает Task 42' на свободном Executor'е
Speculative Execution и Data Skew: почему не помогает¶
При Data Skew медленная Task медленная не потому что Executor плохой - а потому что данных больше. Дублирующая Task получит те же самые данные через shuffle read и будет столь же медленной. Результат: вдвое больше CPU и памяти потрачено впустую, Stage всё равно ждёт 15 минут.
Для Data Skew нужны другие инструменты:
- AQE Skew Join (автоматически)
- Salting ключей группировки
repartition()с увеличением числа партиций
Влияние на ресурсы кластера¶
Трейдофф: каждая speculative copy потребляет дополнительный CPU-слот, память на Executor'е и сетевой трафик (shuffle read дважды). При слишком агрессивных настройках можно перегрузить кластер дублирующими Task'ями.
Признаки слишком агрессивной спекуляции:
- В Spark UI видно много
KILLEDTasks (убитые дубли) numSpeculativeTasksблизко к общему числу Tasks- Утилизация кластера выросла, но Stage завершается не быстрее
Идемпотентность: когда speculation опасен¶
Почему side effects - проблема¶
Speculative copy - это полное повторное выполнение Task с нуля на другом Executor'е. Если Task имеет побочные эффекты (запись в БД, вызов API, отправка сообщений), они выполнятся дважды.
Почему HDFS/S3 безопасны¶
При записи в HDFS или S3 Spark использует механизм Output Committer:
- Каждая Task пишет во временный файл (
_temporary/0/task_attempt_XXXXX/) - После успешного завершения Task - атомарный rename в финальный путь
- Если оригинальная и speculative Task завершились почти одновременно, одна из них перезапишет файл другой - результат идентичен, дублирования нет
- При использовании S3 нужен S3A Committer (алгоритм
stagingилиmagic) - стандартный rename на S3 не атомарен
# Для S3 + speculation: обязательно использовать staging committer
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "staging")
spark.conf.set("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol")
Операции, где нужно отключить speculation¶
# Запись в JDBC (PostgreSQL, MySQL, etc.)
df.write.jdbc(url, table, mode="append")
# Каждая Task делает batch INSERT → дублирование строк при speculation
# Решение: выключить speculation или использовать upsert-режим
# Вызов внешнего API из UDF
@udf("string")
def notify_api(order_id):
import requests
requests.post("https://api.example.com/notify", json={"order_id": order_id})
return "sent"
# При speculation: уведомление отправится дважды
# Kafka sink (без idempotent producer)
df.write.format("kafka") \
.option("kafka.bootstrap.servers", "...") \
.save()
# По умолчанию: at-least-once delivery → дублирование при speculation
Speculative Execution и облачная инфраструктура¶
Spot / Preemptible инстансы¶
В облачных средах (AWS Spot, GCP Preemptible, Azure Spot) инстансы могут быть отозваны в любой момент. Это выглядит как случайный медленный или упавший Executor.
Для кластеров на Spot-инстансах speculation особенно ценен: он позволяет завершить Task до того, как инстанс будет полностью отозван и Executor упадёт. Если speculation успел завершить копию - Stage продолжается без перезапуска.
# Для Spot-кластеров: более агрессивные настройки
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "2.0") # ждём больше, т.к. spot медленнее
spark.conf.set("spark.speculation.quantile", "0.5") # начинаем раньше
Heterogeneous кластер¶
В облаке кластер может содержать инстансы разных типов (разная скорость CPU, IO). Speculation компенсирует неоднородность: медленная Task на слабом инстансе получит копию на быстром.
Диагностика в Spark UI¶
Вкладка Stages → Tasks¶
Что видно в Tasks-таблице:
Task ID | Attempt | Status | Duration | Executor | Speculative
42 | 0 | KILLED | 14m | worker-1 | No ← убит как страглер
42 | 1 | SUCCESS | 2m | worker-5 | Yes ← speculative copy победила
67 | 0 | SUCCESS | 28s | worker-3 | No ← обычный Task
67 | 1 | KILLED | 25s | worker-7 | Yes ← оригинал победил первым
Timeline View¶
В Stage → Event Timeline видно "длинный хвост" - одна полоска Task, которая тянется далеко вправо после того, как остальные уже закрасились зелёным. Момент появления второй полоски (speculative copy) на другом Executor'е обозначает запуск спекуляции.
Программный мониторинг¶
# SparkMeasure - сторонняя библиотека для детальных метрик
from sparkmeasure import StageMetrics
sm = StageMetrics(spark)
sm.begin()
df.groupBy("region").count().collect()
sm.end()
# В отчёте ищем: numSpeculativeTasks
sm.print_report()
# Альтернатива: через SparkContext.statusTracker
tracker = spark.sparkContext._jsc.sc().statusTracker()
for stage_id in tracker.getActiveStageIds():
info = tracker.getStageInfo(stage_id)
if info.isDefined():
print(f"Stage {stage_id}: {info.get().numActiveTasks()} tasks")
Взаимодействие с AQE¶
AQE (Adaptive Query Execution) и Speculative Execution решают разные проблемы:
Важно: AQE и Speculation можно включать одновременно. AQE уменьшает риск Data Skew страглеров (убирает основную причину), а Speculation справляется с оставшимися инфраструктурными проблемами.
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.speculation", "true")
# Вместе: AQE снижает вероятность straggler из-за skew,
# Speculation ловит оставшиеся случайные задержки
Практика¶
1. Симуляция страглера и наблюдение speculation¶
import time
from pyspark.sql.functions import udf
from pyspark.sql.types import LongType
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "2.0")
spark.conf.set("spark.speculation.quantile", "0.5")
# UDF, которая искусственно тормозит на определённом executor
@udf(LongType())
def slow_udf(partition_id):
import socket
import time
# Искусственно замедляем первые партиции
if partition_id % 10 == 0:
time.sleep(30) # Тормозим каждую 10-ю партицию
return partition_id * 2
df = spark.range(0, 1000, 1, numPartitions=100) \
.selectExpr("id", "id % 100 AS partition_id") \
.withColumn("result", slow_udf("partition_id"))
# Пока выполняется - открыть Spark UI http://localhost:4040
# → Jobs → Stage → Tasks
# Наблюдать Tasks с Speculative=Yes
df.count()
2. Сравнение с включённой и выключенной спекуляцией¶
import time
# Тест 1: без спекуляции
spark.conf.set("spark.speculation", "false")
t0 = time.time()
df.count()
t_no_spec = time.time() - t0
print(f"Без speculation: {t_no_spec:.1f}s")
# Тест 2: со спекуляцией
spark.conf.set("spark.speculation", "true")
t0 = time.time()
df.count()
t_spec = time.time() - t0
print(f"Со speculation: {t_spec:.1f}s")
print(f"Ускорение: {t_no_spec/t_spec:.1f}×")
# При искусственном страглере: ускорение = время страглера / время спекуляции
3. Проверка идемпотентности перед включением¶
# Перед включением speculation для Job с write-операциями,
# проверить: является ли операция идемпотентной?
# Безопасно (overwrite mode - перезапись):
df.write.mode("overwrite").parquet("/output/") # ✅
# Безопасно (upsert по ключу):
df.write.format("delta").option("mergeSchema", "true").save("/delta_table/") # ✅ (Delta ACID)
# ОПАСНО (append в JDBC без ключа):
df.write.mode("append").jdbc(url, "my_table") # ❌ дублирование при speculation
# ОПАСНО (side effect UDF):
@udf("string")
def send_email(user_id):
# Отправка email - не идемпотентна
email_service.send(user_id, "Welcome!")
return "sent"
# При speculation: email придёт дважды ❌
4. Тюнинг для production¶
# HDFS batch ETL: стандартные настройки
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "1.5")
spark.conf.set("spark.speculation.quantile", "0.75")
# Cloud-кластер на Spot-инстансах: более агрессивно
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "1.5")
spark.conf.set("spark.speculation.quantile", "0.5") # начинаем раньше
# Кластер с известным skew (AQE снимает часть проблем):
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "3.0") # высокий порог, чтобы не дублировать skew
Когда включать, когда отключать¶
| Сценарий | Speculation | Примечание |
|---|---|---|
| Batch ETL в HDFS/S3 (overwrite) | Включить | Безопасно, помогает при аппаратных проблемах |
| Batch ETL в HDFS/S3 (append) | Включить | S3A Committer гарантирует идемпотентность |
| Запись в Delta Lake / Iceberg | Включить | ACID-движки гарантируют exactly-once |
| JDBC append без ключа | Выключить | Дублирование строк |
| Kafka sink (at-least-once) | Выключить | Дублирование сообщений |
| REST API вызовы из UDF | Выключить | Двойные побочные эффекты |
| ML training (без write) | Включить | Нет side effects |
| Structured Streaming | Выключить | Micro-batch semantics несовместимы |
| Spot / Preemptible кластер | Включить | Компенсирует нестабильность инстансов |
| Кластер с сильным Data Skew | Выключить / AQE | Speculation не поможет, тратит ресурсы |
Резюме¶
Speculative Execution - защита от случайных аппаратных страглеров. Алгоритм:
- Ждать пока завершится
quantile(75%) Tasks Stage - Вычислить медиану времени выполнения
- Запустить дублирующую Task'у для тех, кто медленнее
median × multiplier(1.5×) - Кто первый завершился - тот победил, второй убит
Помогает при аппаратных проблемах (медленный диск, перегрузка узла, Spot-отзыв) и неоднородных кластерах.
Не помогает при Data Skew - дублирующая Task получит те же данные и будет такой же медленной.
Опасен при non-idempotent операциях: JDBC INSERT, REST API, Kafka at-least-once - произойдёт двойное выполнение. Перед включением убедиться в идемпотентности всех операций записи.