İçeriğe atla

Dead Letter Queue Stratejileri: Dayanıklı Olay-Güdümlü Sistemler için Production-Ready Kalıplar

Event-driven sistemler için production-ready DLQ kalıpları: monitoring, circuit breaker, exponential backoff, recovery ve kaçınılması gereken anti-pattern'lar.

Ayhan Sipahi Ayhan Sipahi

Dead Letter Queue, bir tüketicinin retry bütçesini tükettikten sonra işleyemediği mesajları tutan kuyruktur. DLQ olmadan bir poison pill ya birincil kuyruğu baştan (head-of-line) tıkar ya da başarısız handler ile birlikte sessizce kaybolur; her iki sonuçta da hem event hem de bir şeylerin yanlış gittiğine dair operasyonel sinyal kaybedilir. DLQ, “işlenecek mesajlar” ile “insan ya da tool müdahalesi gerektiren mesajlar” arasındaki görev ayrımıdır ve yalnızca etrafındaki retry politikası, uyarı ve replay araçları birlikte tasarlandığında işe yarar.

İşe yarayan varsayılan şu: retry’ı transport seviyesinde sınırlamak, hataları mesaj sınıfına göre yönlendirmek ve DLQ derinliği yerine mesaj yaşına alarm kurmak. Derinlik kaç mesajın başarısız olduğunu söyler. Yaş ise kimsenin ne kadar süredir bakmadığını. Geri kalan her şey (jitter, circuit breaker, replay araçları) bu varsayılanı yük altında ayakta tutmak içindir.

DLQ’nun Rolü#

DLQ, başarıyla işlenemeyen mesajlar için güvenlik ağınızdır. DLQ yönetimi olmadan başarısız mesajlar şu üç şeyden birini yapar:

  1. Sonsuza kadar kaybolur (sessiz hatalar)
  2. Tüm kuyruğu bloke eder (poison pill problemi)
  3. Sonsuz retry döngüleri oluşturur (cascade hatalar)

Kuyruğu faydalı kılan iki şey var: her mesaja iliştirilen hata bağlamı ve ana kuyruğa geri dönüş için yazılı bir yol.

DLQ Implementation Pattern’ları#

Pattern 1: Jitter ile Exponential Backoff#

En yaygın pattern. İşe yarayıp yaramadığını iki ayrıntı belirler: jitter ve vazgeçmeden önce mesaja iliştirdiğiniz retry bağlamı.

class ResilientMessageProcessor {
  async processWithBackoff(message: Message, maxRetries = 5) {
    let retryCount = 0;
    let lastError;

    while (retryCount < maxRetries) {
      try {
        return await this.process(message);
      } catch (error) {
        lastError = error;
        retryCount++;

        // Thundering herd'i önlemek için jitter ekle
        const baseDelay = Math.pow(2, retryCount - 1) * 1000;
        const jitter = Math.random() * 1000;
        const delay = baseDelay + jitter;

        await this.sleep(delay);

        // Retry bağlamıyla mesajı zenginleştir
        message.metadata = {
          ...message.metadata,
          retryCount,
          lastError: error.message,
          retryTimestamp: new Date().toISOString(),
          backoffDelay: delay
        };
      }
    }

    // Max retry aşıldı - tam bağlamla DLQ'ya gönder
    await this.sendToDLQ(message, lastError, retryCount);
  }

  async sendToDLQ(message: Message, error: Error, attempts: number) {
    const dlqPayload = {
      originalMessage: message,
      failureReason: {
        errorMessage: error.message,
        errorStack: error.stack,
        errorType: error.constructor.name,
        timestamp: new Date().toISOString()
      },
      processingContext: {
        totalAttempts: attempts,
        firstAttempt: message.metadata?.firstAttempt || new Date().toISOString(),
        finalAttempt: new Date().toISOString(),
        processingDuration: this.calculateProcessingTime(message)
      },
      environmentContext: {
        nodeVersion: process.version,
        hostname: os.hostname(),
        memoryUsage: process.memoryUsage()
      }
    };

    await this.dlqClient.send(dlqPayload);

    // DLQ metriklerini artır
    this.metrics.dlqMessages.inc({
      errorType: error.constructor.name,
      messageType: message.type
    });
  }
}

Pattern 2: Circuit Breaker DLQ#

Downstream servis hataları için:

class CircuitBreakerDLQ {
  private failures = new Map<string, { count: number, lastFailure: Date }>();
  private circuitState: 'CLOSED' | 'OPEN' | 'HALF_OPEN' = 'CLOSED';

  async processMessage(message: Message) {
    const serviceKey = this.extractServiceKey(message);

    if (this.isCircuitOpen(serviceKey)) {
      // Deneme bile yapma - direkt DLQ'ya circuit breaker sebebiyle
      return this.sendToDLQ(message, new Error('Circuit breaker open'), {
        circuitState: this.circuitState,
        failureCount: this.failures.get(serviceKey)?.count || 0
      });
    }

    try {
      const result = await this.processWithTimeout(message, 30000);
      this.recordSuccess(serviceKey);
      return result;
    } catch (error) {
      this.recordFailure(serviceKey);

      if (this.shouldOpenCircuit(serviceKey)) {
        this.openCircuit(serviceKey);
      }

      throw error; // Normal retry mantığının halletmesine bırak
    }
  }

  private isCircuitOpen(serviceKey: string): boolean {
    const failure = this.failures.get(serviceKey);
    if (!failure) return false;

    // Son 5 dakikada 5+ hata varsa devreyi aç (eşikler yapılandırılabilir)
    return failure.count >= 5 &&
           (Date.now() - failure.lastFailure.getTime()) < 300000;
  }
}

Pattern 3: Content-Based DLQ Routing#

Farklı mesaj tipleri farklı DLQ stratejileri gerektirir:

class SmartDLQRouter {
  private dlqStrategies = new Map([
    ['payment', { maxRetries: 10, alertLevel: 'CRITICAL' }],
    ['notification', { maxRetries: 3, alertLevel: 'WARNING' }],
    ['analytics', { maxRetries: 1, alertLevel: 'INFO' }],
  ]);

  async processMessage(message: Message) {
    const messageType = message.headers?.type || 'default';
    const strategy = this.dlqStrategies.get(messageType) || { maxRetries: 3, alertLevel: 'WARNING' };

    try {
      return await this.processWithStrategy(message, strategy);
    } catch (error) {
      // Mesaj tipi ve hataya göre uygun DLQ'ya yönlendir
      const dlqTopic = this.selectDLQTopic(messageType, error);
      await this.sendToSpecificDLQ(dlqTopic, message, error, strategy);
    }
  }

  private selectDLQTopic(messageType: string, error: Error): string {
    // Kritik mesajlar yüksek öncelikli DLQ'ya gider
    if (messageType === 'payment') {
      return 'payment-dlq-critical';
    }

    // Geçici hatalar retry DLQ'ya gider
    if (this.isTemporaryError(error)) {
      return 'retry-dlq';
    }

    // Kalıcı hatalar inceleme DLQ'sına gider
    return 'investigation-dlq';
  }
}

DLQ Monitoring: Kapsamlı Metrikler#

Tek başına derinlik yalnızca kuyruğun dolduğunu söyler. Yanında şunları toplayın:

class DLQMonitoring {
  private metrics = {
    // Temel metrikler
    dlqDepth: new Gauge('dlq_depth'),
    dlqRate: new Counter('dlq_messages_total'),

    // Gelişmiş metrikler
    dlqMessageAge: new Histogram('dlq_message_age_seconds'),
    errorPatterns: new Counter('dlq_error_patterns', ['error_type', 'message_type']),
    retrySuccessRate: new Gauge('dlq_retry_success_rate'),

    // Business metrikler
    revenueImpact: new Gauge('dlq_revenue_impact_dollars'),
    customerImpact: new Counter('dlq_customer_impact', ['severity'])
  };

  async trackDLQMessage(message: DLQMessage) {
    // Hata kalıplarını takip et
    this.metrics.errorPatterns.inc({
      error_type: message.failureReason.errorType,
      message_type: message.originalMessage.type
    });

    // Business impact hesapla
    const impact = await this.calculateBusinessImpact(message);
    this.metrics.revenueImpact.set(impact.revenue);
    this.metrics.customerImpact.inc({ severity: impact.severity });

    // Mesaj yaşı takibi
    const messageAge = Date.now() - new Date(message.originalMessage.timestamp).getTime();
    this.metrics.dlqMessageAge.observe(messageAge / 1000);
  }
}

DLQ Recovery Stratejileri#

Strateji 1: Bilinen Hatalar için Otomatik Recovery#

Yalnızca daha önce teşhis ettiğiniz hata imzalarını otomatikleştirin. Tanınmayan her şey insanı bekler; çünkü yanlış bir otomatik düzeltme, bozuk veriyi makine hızında yeniden oynatır.

class KnownIssueDLQRecovery {
  // Her kayıt bir hata imzasını deterministik bir düzeltmeye (Fix implementasyonu) bağlar
  private fixes = new Map<string, Fix>([
    ['SchemaValidationError:missing_currency', fillDefaultCurrency],
    ['UpstreamTimeout', replayUnchanged],
  ]);

  async analyzeAndRecover() {
    const dlqMessages = await this.fetchDLQMessages();

    // Hata imzasına göre grupla ki tek karar tüm batch'i kapsasın
    const errorGroups = this.groupByErrorPattern(dlqMessages);

    for (const [pattern, messages] of errorGroups.entries()) {
      const fix = this.fixes.get(pattern);

      if (fix) {
        await this.applyKnownFix(messages, fix);
      } else {
        // Bilinmeyen imza: replay'den önce bir insan sınıflandırır
        await this.createTicket(pattern, messages);
      }
    }
  }

  private async applyKnownFix(messages: DLQMessage[], fix: Fix) {
    for (const message of messages) {
      try {
        const fixedMessage = await fix.apply(message);
        await this.mainQueue.send(fixedMessage);

        // Yalnızca ana kuyruk düzeltilmiş mesajı kabul ettikten sonra sil
        await this.dlq.delete(message);
      } catch (error) {
        // DLQ'da bırak; başarısız bir düzeltme kanıtı yok etmemeli
        await this.recordFailedRepair(message, error);
      }
    }
  }
}

Strateji 2: Progressive Recovery#

class ProgressiveDLQRecovery {
  async recoverInWaves(batchSize = 10) {
    let recovered = 0;
    let failed = 0;
    let wave = 0;

    while (true) {
      const batch = await this.dlq.receiveMessages({ MaxMessages: batchSize });
      if (batch.length === 0) break;

      // Batch'leri aralarında exponential gecikmeyle işle
      const results = await this.processBatch(batch);

      recovered += results.successful;
      failed += results.failed;
      wave++;

      // Hata oranı yüksekse duraklat ve uyar
      const failureRate = failed / (recovered + failed);
      if (failureRate > 0.5) {
        await this.alertOncallTeam(
          `DLQ recovery failure rate: ${(failureRate * 100).toFixed(0)}%`
        );
        await this.sleep(60000); // 1 dakika bekle
      }

      // Toplam hata sayısına göre değil, dalga sayısına göre geri çekil
      await this.sleep(Math.min(1000 * Math.pow(2, wave), 30000));
    }
  }
}

Cloud Provider DLQ Özellikleri#

AWS SQS DLQ#

# CloudFormation şablonu
Resources:
  MainQueue:
    Type: AWS::SQS::Queue
    Properties:
      RedrivePolicy:
        deadLetterTargetArn: !GetAtt DLQ.Arn
        maxReceiveCount: 3
      MessageRetentionPeriod: 1209600  # 14 gün

  DLQ:
    Type: AWS::SQS::Queue
    Properties:
      MessageRetentionPeriod: 1209600  # 14 gün

  DLQAlarm:
    Type: AWS::CloudWatch::Alarm
    Properties:
      AlarmName: DLQ-HighDepth
      MetricName: ApproximateNumberOfMessagesVisible
      Namespace: AWS/SQS
      Dimensions:
        - Name: QueueName
          Value: !GetAtt DLQ.QueueName
      Statistic: Average
      Threshold: 10
      ComparisonOperator: GreaterThanThreshold

Azure Service Bus DLQ#

// Otomatik DLQ yönetimi
var options = new ServiceBusProcessorOptions
{
    MaxConcurrentCalls = 10,
    MaxAutoLockRenewalDuration = TimeSpan.FromMinutes(10),
    // Mesajlar MaxDeliveryCount sonrası otomatik DLQ'ya gider (varsayılan: 10)
    SubQueue = SubQueue.None  // Ana kuyruk
};

// Recovery için DLQ'ya eriş
var dlqProcessor = client.CreateProcessor(
    queueName,
    new ServiceBusProcessorOptions { SubQueue = SubQueue.DeadLetter }
);

GCP Pub/Sub DLQ#

# Terraform yapılandırması
resource "google_pubsub_subscription" "main" {
  name  = "main-subscription"
  topic = google_pubsub_topic.main.name

  dead_letter_policy {
    dead_letter_topic  = google_pubsub_topic.dlq.id
    max_delivery_attempts = 5
  }

  retry_policy {
    minimum_backoff = "10s"
    maximum_backoff = "600s"
  }
}

Kaçınılması Gereken DLQ Anti-Pattern’ları#

  1. “Kur ve Unut” Anti-Pattern’ı

    • Monitoring olmadan DLQ oluşturmak
    • DLQ’daki mesajları hiç işlememek
    • DLQ derinliği için alarm kurmamak
  2. “Sonsuz Retry” Anti-Pattern’ı

    • Maksimum retry sınırının olmaması
    • Tüm hata tipleri için aynı retry gecikmesi
    • Downstream hatalar için circuit breaker bulunmaması
  3. “Kara Delik” Anti-Pattern’ı

    • Bağlamsız DLQ mesajları
    • Hata sınıflandırmasının olmaması
    • Recovery prosedürlerinin tanımsız olması

Production DLQ Checklist#

  • Uygun retention periyodları yapılandır (minimum 14 gün)
  • DLQ derinlik alarmları kur (> 10 mesaj)
  • DLQ yaş metriklerini monitor et (1 saatten eski mesajlar)
  • Bilinen hata kalıpları için otomatik recovery uygula
  • Manuel araştırma için runbook’lar oluştur
  • DLQ mesajlarından business impact metriklerini takip et
  • Ekip standup’larında düzenli DLQ gözden geçirmeleri
  • Yüksek hata oranları sırasında DLQ davranışını load test et

Yaygın DLQ Hata Kalıpları#

Sessiz Payment Başarısızlığı#

DLQ’lar izlenmediğinde payment’lar günlerce sessizce başarısız olabilir. Mesajlar alarm üretmeden birikir ve ilk sinyal bir dashboard yerine müşteri şikâyeti olur. Çözüm: Yalnızca ana kuyruk metriklerine değil, DLQ yaşına da alarm kurun.

Recovery Sırasında Thundering Herd#

Downstream servis kesintisi sırasında, jitter olmadan eş zamanlı retry girişimleri toparlanmaya çalışan servisi aşırı yükler ve kesinti süresini uzatır. Çözüm: Retry girişimlerini yaymak için exponential backoff’a her zaman jitter ekleyin.

Poison Pill Bloklaması#

Hatalı biçimlendirilmiş bir mesaj her denemede consumer servisi çökertebilir. Doğru DLQ routing olmadan yüksek trafik dönemlerinde sonraki tüm mesajları bloke eder. Çözüm: Circuit breaker’lar ve farklı hata tipleri için ayrı DLQ’lar uygulayın.

Sonuç#

Varsayılan, çoğu event-driven iş yükü için geçerlidir: tüketici başına bir DLQ, transport seviyesinde sınırlanmış retry ve mesaj yaşına kurulmuş bir alarm. Mesaj sınıfı ekonomiyi değiştirdiğinde bu varsayılanı esnetin. Payment event’leri kendi DLQ’sunu, daha sıkı bir alarmı ve yazılı bir replay prosedürünü hak eder. Analytics event’leri çoğunlukla etmez; paylaşılan tek bir DLQ’yu düzenli aralıklarla gözden geçirmek yeterlidir.

Bir sonraki deploy’dan önce DLQ’daki bir mesajı elle replay edin. Bu iş alarmı okumaktan uzun sürüyorsa eksik olan araçlardır.


İlgili Okuma: Olay-güdümlü sistem araçları ve kalıplarının daha geniş bir genel bakışı için olay-güdümlü mimari araçları kapsamlı rehberini görün.

Kaynaklar#

İlgili yazılar