Speculative Execution: straggler detection и когда отключать

Spark запускает дублирующую Task для медленных стрэглеров. Разбираем механизм обнаружения, конфигурацию, идемпотентность и случаи, когда speculation делает хуже.

core optimization

Проблема хвостовой латентности в распределённых системах

В 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 видно много KILLED Tasks (убитые дубли)
  • 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 - защита от случайных аппаратных страглеров. Алгоритм:

  1. Ждать пока завершится quantile (75%) Tasks Stage
  2. Вычислить медиану времени выполнения
  3. Запустить дублирующую Task'у для тех, кто медленнее median × multiplier (1.5×)
  4. Кто первый завершился - тот победил, второй убит

Помогает при аппаратных проблемах (медленный диск, перегрузка узла, Spot-отзыв) и неоднородных кластерах.

Не помогает при Data Skew - дублирующая Task получит те же данные и будет такой же медленной.

Опасен при non-idempotent операциях: JDBC INSERT, REST API, Kafka at-least-once - произойдёт двойное выполнение. Перед включением убедиться в идемпотентности всех операций записи.