İçeriğe atla

Kafka vs SQS vs EventBridge: Hangisini Seçmeli?

Olay güdümlü araçlara derin bakış: Kafka, SQS ve EventBridge, mesaj teslimat kalıpları, DLQ stratejileri ve bunların AWS, Azure, GCP karşılıkları.

Ayhan Sipahi Ayhan Sipahi

Olay güdümlü araç seçimi hype’tan çok üç şekilden birine oturtmakla ilgili: kuyruk, fan-out topic ya da router. O şekil için bulutunuzun zaten sunduğu yönetilen servisle başlayın. Kafka’ya ancak event replay, partition bazlı sıralama ya da yönetilen kuyruk kotasının karşılamadığı bir hacim gerektiğinde geçin.

Kararın geri kalanı teslimat garantilerine, mesaj boyutu limitlerine ve kimsenin işleyemediği bir mesaja ne olacağına bakar.

Mesaj Kalıpları: Üç Temel Şekil#

Her pattern farklı araçlar ve trade-off’lar getirir:

1’e 1 (Kuyruk Kalıbı)#

  • Mesaj tek consumer tarafından tüketilir (FIFO ile tam olarak bir kez)
  • Kullanım: Task işleme, iş dağıtımı, async job’lar
  • Araçlar: SQS, Azure Service Bus Queues, Cloud Tasks

1’e Çok (Topic/Fan-out Kalıbı)#

  • Mesaj birden fazla subscriber’a iletilir; her biri kendi kopyasını alır
  • Kullanım: Event yayınlama, bildirimler, audit log
  • Araçlar: SNS, Azure Service Bus Topics, Cloud Pub/Sub

Çok’a Çok (Event Mesh)#

  • Birden fazla producer/consumer arası karmaşık routing, filtreleme
  • Kullanım: Microservices iletişimi, domain event’leri
  • Araçlar: EventBridge, Azure Event Grid, Eventarc

Araç Kategorileri#

Basit Kuyruk Servisleri#

AWS SQS (Simple Queue Service)#

Güçlü olduğu yer: Çok basit kuyruk operasyonları, serverless entegrasyon, otomatik ölçekleme

Alma ve DLQ konfigürasyonu:

// Long polling boş receive'leri ucuza getirir
const params = {
  QueueUrl: 'https://sqs.us-east-1.amazonaws.com/123/my-queue',
  ReceiveMessageWaitTimeSeconds: 20,  // Long polling
  MaxNumberOfMessages: 10,
  VisibilityTimeout: 30,  // İşleme penceresi
  MessageAttributeNames: ['All']
};

// DLQ'nun kendisine sadece uzun retention gerekir
const dlqParams = {
  QueueName: 'my-queue-dlq',
  Attributes: {
    MessageRetentionPeriod: '1209600'  // 14 gün, SQS'in üst sınırı
  }
};

// Redrive policy DLQ'ya değil, KAYNAK kuyruğa yazılır
const sourceQueueParams = {
  QueueName: 'my-queue',
  Attributes: {
    RedrivePolicy: JSON.stringify({
      deadLetterTargetArn: dlqArn,  // my-queue-dlq'nun ARN'i
      maxReceiveCount: 3  // 3 teslimattan sonra DLQ'ya taşı
    })
  }
};

Teslimat garantileri:

  • Standard Queue: En az bir kez (olası duplikasyonlar)
  • FIFO Queue: Tam olarak bir kez işleme
  • Mesaj sıralaması: Sadece FIFO
  • Max mesaj boyutu: 1MB (Ağustos 2025’te 256KB’den yükseltildi)

Note

Mesaj boyutu limitindeki bu 4x artış, daha büyük veri alışverişi gerektiren AI, IoT ve karmaşık uygulama entegrasyon workload’ları için faydalıdır. AWS Lambda’nın event source mapping’i de yeni 1MB payload’ları destekleyecek şekilde güncellendi.

SQS ne zaman doğru seçim:

  • Microservices’i decouple etme
  • Batch job işleme
  • Serverless mimariler (Lambda trigger’ları)
  • Basit task kuyrukları

Azure Service Bus Queues#

Enterprise özellikleriyle SQS’in Azure eşleniği:

// Session'lar ve DLQ handling ile Service Bus
var client = new ServiceBusClient(connectionString);
var processor = client.CreateProcessor(queueName, new ServiceBusProcessorOptions
{
    MaxConcurrentCalls = 10,
    AutoCompleteMessages = false,
    MaxAutoLockRenewalDuration = TimeSpan.FromMinutes(5),
    SubQueue = SubQueue.DeadLetter  // DLQ'ya erişim
});

// Duplikasyon algılamalı mesaj
var message = new ServiceBusMessage(body)
{
    MessageId = Guid.NewGuid().ToString(),  // Deduplication için
    SessionId = sessionId,  // Sıralı işleme için
    TimeToLive = TimeSpan.FromMinutes(5)
};

SQS’ten temel farklar:

  • Sıralı işleme için built-in session’lar
  • Duplikasyon algılama (yapılandırılabilir pencere)
  • Zamanlanmış mesajlar
  • Mesaj boyutu: 256KB (standard), 100MB (premium)

Google Cloud Tasks#

HTTP target entegrasyonlu GCP task kuyruğu:

import { CloudTasksClient } from '@google-cloud/tasks';

const client = new CloudTasksClient();
const parent = client.queuePath(project, location, queue);

const task = {
    httpRequest: {
        httpMethod: 'POST',
        url: 'https://example.com/process',
        headers: { 'Content-Type': 'application/json' },
        body: Buffer.from(JSON.stringify(payload))
    },
    scheduleTime: { seconds: Math.floor(timestamp / 1000) }
};

const response = await client.createTask({ parent, task });

Pub/Sub Sistemleri#

AWS SNS (Simple Notification Service)#

1’e çok mesaj dağıtımı:

// Akıllı routing için filter policy'li SNS
const publishParams = {
  TopicArn: 'arn:aws:sns:us-east-1:123:my-topic',
  Message: JSON.stringify(event),
  MessageAttributes: {
    eventType: { DataType: 'String', StringValue: 'ORDER_CREATED' },
    priority: { DataType: 'Number', StringValue: '1' }
  }
};

// Filter'lı subscription
const subscriptionPolicy = {
  eventType: ['ORDER_CREATED', 'ORDER_UPDATED'],
  priority: [{ numeric: ['>', 0] }]
};

SNS + SQS Pattern (Fanout):

Producer

SNS Topic

SQS Queue 1

SQS Queue 2

SQS Queue 3

Consumer 1

Consumer 2

Consumer 3

Teslimat garantileri:

  • En az bir kez teslimat
  • Mesaj sıralaması yok
  • Exponential backoff ile retry
  • Başarısız teslimatlar için DLQ desteği

Azure Service Bus Topics#

SNS’ten daha sofistike:

// Birden fazla subscription ve filter'lı topic
var adminClient = new ServiceBusAdministrationClient(connectionString);

await adminClient.CreateSubscriptionAsync(
    new CreateSubscriptionOptions(topicName, subscriptionName),
    new CreateRuleOptions("OrderFilter",
        new SqlRuleFilter("EventType = 'OrderCreated' AND Priority > 5"))
);

Gelişmiş özellikler:

  • SQL benzeri filtre kuralları
  • Sıralama için mesaj session’ları
  • Duplikasyon algılama
  • Sebep takipli dead-lettering

Google Cloud Pub/Sub#

Global mesaj dağıtımı:

import { PubSub } from '@google-cloud/pubsub';

const pubsub = new PubSub();
// Publisher açıkça izin vermezse ordering key yok sayılır
const topic = pubsub.topic(topicId, { enableMessageOrdering: true });

// Ordering key ile yayınlama
const messageId = await topic.publishMessage({
    data: Buffer.from(data),
    orderingKey: 'user-123',
    attributes: {
        event_type: 'user_updated',
        version: '2'
    }
});

Event Routing Servisleri#

AWS EventBridge#

Kural tabanlı event routing:

// İçerik tabanlı routing ile EventBridge
const rule = {
  Name: 'OrderProcessingRule',
  EventPattern: JSON.stringify({
    source: ['order.service'],
    'detail-type': ['Order Created'],
    detail: {
      amount: [{ numeric: ['>', 100] }],
      country: ['US', 'UK', 'DE']
    }
  }),
  Targets: [
    {
      Arn: lambdaArn,
      RetryPolicy: {
        MaximumRetryAttempts: 2,
        MaximumEventAgeInSeconds: 3600
      },
      DeadLetterConfig: {
        Arn: dlqArn
      }
    }
  ]
};

Hesaplar arası event paylaşımı:

// PutPermission, event'i alan hesabın bus'ında çalışır.
// Ya tek bir hesap adı verin ya da '*' ile organizasyon koşulunu kullanın.
const eventBusPolicy = {
  EventBusName: 'default',
  StatementId: 'AllowOrgAccess',
  Action: 'events:PutEvents',
  Principal: '*',
  Condition: {
    Type: 'StringEquals',
    Key: 'aws:PrincipalOrgID',
    Value: 'o-1234567890'  // PutPermission'ın kabul ettiği tek koşul anahtarı
  }
};

// detail-type filtresi burada değil, alıcı bus'ın kurallarında yapılır

Azure Event Grid#

Güçlü filtreleme ile Azure eşleniği:

{
  "filter": {
    "includedEventTypes": ["Microsoft.Storage.BlobCreated"],
    "subjectBeginsWith": "/blobServices/default/containers/images/",
    "advancedFilters": [
      {
        "operatorType": "NumberGreaterThan",
        "key": "data.contentLength",
        "value": 1048576
      }
    ]
  }
}

Google Cloud Eventarc#

GCP’nin birleşik eventing’i:

# Eventarc trigger konfigürasyonu
apiVersion: eventarc.cnrm.cloud.google.com/v1beta1
kind: EventarcTrigger
metadata:
  name: storage-trigger
spec:
  location: us-central1
  matchingCriteria:
  - attribute: type
    value: google.cloud.storage.object.v1.finalized
  - attribute: bucket
    value: my-bucket
  destination:
    cloudRunService:
      name: process-image
      region: us-central1

Stream Processing Platformları#

Apache Kafka#

Teslimat semantiği yapılandırılabilen açık kaynak event streaming:

// Gerçek zamanlı işleme için Kafka Streams
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-processor");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, "exactly_once_v2");
props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);

KStream<String, Order> orders = builder.stream("orders");
KTable<String, Long> orderCounts = orders
    .filter((k, v) -> v.getAmount() > 100)
    .groupByKey()
    .count(Materialized.as("order-counts-store"));

// Kafka Streams ile DLQ handling
orders.foreach((key, value) -> {
    try {
        processOrder(value);
    } catch (Exception e) {
        producer.send(new ProducerRecord<>("orders-dlq", key, value));
    }
});

Kafka teslimat semantiği:

  • En fazla bir kez: Gönder ve unut (acks=0)
  • En az bir kez: Varsayılan (acks=1 veya all)
  • Tam olarak bir kez: Transaction’larla (enable.idempotence=true)

Cloud Streaming Eşlenikleri#

AWS Kinesis Data Streams#

import {
  KinesisClient,
  RegisterStreamConsumerCommand
} from '@aws-sdk/client-kinesis';

const kinesis = new KinesisClient({ region: 'us-east-1' });

// Enhanced fan-out her consumer'a kendi okuma throughput'unu verir
const consumer = await kinesis.send(new RegisterStreamConsumerCommand({
  StreamARN: streamArn,
  ConsumerName: 'low-latency-consumer'
}));

Azure Event Hubs#

// Kafka protokolü ile Event Hubs
var config = new ConsumerConfig
{
    BootstrapServers = "namespace.servicebus.windows.net:9093",
    SecurityProtocol = SecurityProtocol.SaslSsl,
    SaslMechanism = SaslMechanism.Plain,
    GroupId = "consumer-group"
};

// Uzun süreli depolama için Data Lake'e capture
var captureDescription = new CaptureDescription
{
    Enabled = true,
    IntervalInSeconds = 300,
    SizeLimitInBytes = 314572800,
    Destination = new Destination
    {
        StorageAccountResourceId = "/subscriptions/.../storageAccounts/...",
        BlobContainer = "capture"
    }
};

Google Cloud Dataflow#

# Dataflow, Apache Beam pipeline'larını çalıştırır; yazdığınız SDK Beam'dir
import json
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

options = PipelineOptions(streaming=True, runner='DataflowRunner')

with beam.Pipeline(options=options) as pipeline:
    (pipeline
     | 'Read' >> beam.io.ReadFromPubSub(topic=topic)
     | 'Parse' >> beam.Map(lambda payload: json.loads(payload.decode('utf-8')))
     | 'Window' >> beam.WindowInto(beam.window.FixedWindows(60))
     | 'Filter' >> beam.Filter(lambda item: item['amount'] > 100)
     | 'Write' >> beam.io.WriteToBigQuery(table_spec))

Dead Letter Queue (DLQ) Esasları#

DLQ’lar production dayanıklılığı için kritiktir. Retry’lardan sonra başarıyla işlenemeyen mesajları ele alır.

Temel DLQ kavramları:

  • Başarısız mesajlar için güvenlik ağı
  • Poison pill senaryolarını önler
  • Hata analizi ve kurtarmayı mümkün kılar
  • Kuyruk derinliğinin ötesinde temel monitoring

Kurulum, SQS bölümündeki ikiliyle aynı: DLQ’da uzun retention, kaynak kuyrukta maxReceiveCount. Bu sayıyı retry bütçenize göre ayarlayın. Üç teslimat, geçici hataları soğurur ve zehirli bir mesajı dakikalarca döngüde tutmaz.

Derinlemesine İnceleme: Kapsamlı DLQ stratejileri, monitoring pattern’leri, circuit breaker’lar, ML tabanlı kurtarma ve production dersleri için detaylı rehberimize bakın: Dead Letter Queue Production Stratejileri

Edge ve Hybrid Deployment’lar#

Edge Tarafındaki Kısıtlar#

Edge’deki event sistemleri kendine özgü kısıtlarla çalışır:

// Edge-optimize edilmiş event işleme
class EdgeEventProcessor {
  private localQueue: Queue[] = [];
  private cloudBuffer: Message[] = [];

  async processEvent(event: Event) {
    // Önce lokal işle
    const processed = await this.localProcess(event);

    // Cloud sync için batch'le
    if (this.shouldSyncToCloud(processed)) {
      this.cloudBuffer.push(processed);

      if (this.cloudBuffer.length >= 100 ||
          Date.now() - this.lastSync > 60000) {
        await this.syncToCloud();
      }
    }
  }

  private async syncToCloud() {
    try {
      // Sıkıştır ve batch gönder
      const compressed = this.compress(this.cloudBuffer);
      await this.cloudClient.sendBatch(compressed);
      this.cloudBuffer = [];
      this.lastSync = Date.now();
    } catch (error) {
      // Cloud ulaşılamazsa lokal sakla
      await this.localStorage.store(this.cloudBuffer);
    }
  }
}

Cloudflare Workers ile Kuyruklar#

// Cloudflare Workers Queue Handler
export default {
  async queue(batch: MessageBatch, env: Env): Promise<void> {
    for (const message of batch.messages) {
      try {
        // Edge'de işle
        const result = await processMessage(message.body);

        // Durable Objects veya KV'de sakla
        await env.KV.put(
          `processed:${message.id}`,
          JSON.stringify(result),
          { expirationTtl: 3600 }
        );

        message.ack();
      } catch (error) {
        // Backoff ile retry
        message.retry({ delaySeconds: 30 });
      }
    }
  }
};

AWS IoT Core ile Edge Event’leri#

// Lokal IPC üzerinden IoT Core ile konuşan Greengrass V2 bileşeni
import * as greengrasscoreipc from 'aws-iot-device-sdk-v2/dist/greengrasscoreipc';
import * as model from 'aws-iot-device-sdk-v2/dist/greengrasscoreipc/model';

class EdgeIoTProcessor {
    private ipcClient = greengrasscoreipc.createClient();

    constructor(private deviceId: string) {}

    async connect(): Promise<void> {
        await this.ipcClient.connect();
    }

    async publishEdgeEvent(event: unknown): Promise<void> {
        // Önce cihazda küçült; uplink kapalı olabilir
        const processed = this.processLocally(event);

        const request: model.PublishToIoTCoreRequest = {
            topicName: `edge/${this.deviceId}/events`,
            qos: model.QOS.AT_LEAST_ONCE,
            payload: Buffer.from(JSON.stringify(processed))
        };

        await this.ipcClient.publishToIoTCore(request);
    }

    private processLocally(event: unknown): unknown {
        // Uplink bant genişliği harcamadan filtrele, zenginleştir veya topla
        return event;
    }
}

Cross-Cloud Eşlenikler#

Servis Eşleme Tablosu#

AWSAzureGCPKullanım Alanı
SQSService Bus QueuesCloud TasksBasit kuyruk
SNSService Bus TopicsCloud Pub/SubPub/Sub mesajlaşma
EventBridgeEvent GridEventarcEvent routing
KinesisEvent HubsPub/Sub + DataflowStream processing
Lambda + SQSFunctions + Service BusCloud Run + Pub/SubServerless event’ler
DynamoDB StreamsCosmos DB Change FeedFirestore TriggersDatabase event’leri
Step FunctionsLogic AppsWorkflowsEvent orchestration
MSK (Kafka)Event Hubs (Kafka mode)Confluent CloudKafka-uyumlu

Multi-Cloud Event Bridge Pattern#

// Soyut multi-cloud event interface'i
interface CloudEventAdapter {
  publish(event: CloudEvent): Promise<void>;
  subscribe(handler: EventHandler): Promise<void>;
}

class MultiCloudEventBridge {
  private adapters: Map<string, CloudEventAdapter> = new Map();

  constructor() {
    this.adapters.set('aws', new AWSEventBridgeAdapter());
    this.adapters.set('azure', new AzureEventGridAdapter());
    this.adapters.set('gcp', new GCPEventarcAdapter());
  }

  async publishToAll(event: CloudEvent) {
    const promises = Array.from(this.adapters.values())
      .map(adapter => adapter.publish(event));

    const results = await Promise.allSettled(promises);

    // Kısmi başarısızlıkları ele al
    const failures = results.filter(r => r.status === 'rejected');
    if (failures.length > 0) {
      await this.handleFailures(failures, event);
    }
  }
}

Yetenek Karşılaştırma Matrisi#

Throughput sütunu benchmark sonucu değil, belgelenmiş kotayı ya da sistemin ölçeklendiği ekseni gösterir. Yönetilen kotalar bölgeye göre değişir ve çoğu talep üzerine yükseltilebilir.

AraçThroughput tavanıMesaj BoyutuSıralamaTeslimat GarantisiDLQ Desteği
SQS StandardPratikte sınırsız1MBHayırEn az bir kezEvet
SQS FIFOAPI aksiyonu başına partition başına 300 TPS, batch ile 3K/sn, high-throughput modda daha yüksek1MBEvetTam olarak bir kez işlemeEvet
SNSBölgeye bağlı publish kotası256KBHayırEn az bir kezEvet
KafkaPartition ve broker sayısıyla ölçeklenir1MB defaultPartition başınaYapılandırılabilirManuel
RabbitMQNode sayısı ve kuyruk tipiyle ölçeklenirYapılandırılabilir (max_message_size)OpsiyonelEn az bir kezEvet
EventBridgeBölgeye bağlı PutEvents kotası256KBHayırEn az bir kezEvet
KinesisShard başına 1MB/sn1MBShard başınaEn az bir kezManuel
Azure Service BusMessaging unit sayısıyla ölçeklenir256KB standard, 100MB premiumEvetEn az bir kezEvet
Cloud Pub/SubBölgeye bağlı publish kotası10MBOrdering key başınaEn az bir kezEvet
Redis StreamsInstance boyutuyla ölçeklenirAlan başına 512MBEvetEn az bir kezManuel

Karar Çerçevesi#

Hızlı Karar Ağacı#

Evet

Hayır

AWS

Azure

GCP

Evet

Yüksek

Orta

Event Sistemi Gerekli

Basit Kuyruk?

Cloud Native?

Streaming Gerekli?

SQS

Service Bus

Cloud Tasks

Hacim?

Kafka/MSK

Kinesis/EventHubs

Kullanım Senaryoları#

  • Basit Kuyruklar (SQS/Service Bus): Servisleri decouple etme, iş dağıtımı, serverless işleme
  • Pub/Sub (SNS/Topics): Event yayınlama, fan-out, birden fazla consumer, bildirimler
  • Event Router’lar (EventBridge/EventGrid): Karmaşık routing, multi-service orchestration, SaaS entegrasyonları
  • Streaming (Kafka/Kinesis): Real-time analitik, event sourcing, yönetilen kuyruk kotasını aşan sürekli hacim, event replay

Yaygın Tuzaklar ve Çözümler#

Tuzak 1: Mesaj Boyut Limitleri#

// Çözüm: Claim check pattern
class LargeMessageHandler {
  async send(largePayload: unknown) {
    const body = JSON.stringify(largePayload);

    // 256KB, SNS ve EventBridge tavanı; SQS artık 1MB'a izin veriyor
    if (body.length > 256_000) {
      const s3Key = await this.uploadToS3(largePayload);

      // Payload'ı değil, referansı gönder
      return this.queue.send({
        type: 'large_message',
        s3Key,
        size: body.length
      });
    }

    return this.queue.send(largePayload);
  }
}

Tuzak 2: Zehirli Mesajlar#

// Çözüm: Zehirli mesaj algılama
class PoisonMessageDetector {
  private messageAttempts = new Map<string, number>();

  async process(message: Message) {
    const attempts = this.messageAttempts.get(message.id) || 0;

    if (attempts >= 3) {
      // Zehirli mesaj olarak tanımlandı
      await this.quarantine(message);
      return;
    }

    try {
      await this.processMessage(message);
      this.messageAttempts.delete(message.id);
    } catch (error) {
      this.messageAttempts.set(message.id, attempts + 1);

      if (this.isPoisonPattern(error)) {
        await this.quarantine(message);
      } else {
        throw error; // Retry
      }
    }
  }
}

Tuzak 3: Sıralama Garantileri#

// Çözüm: Partition key stratejisi
class OrderedEventProcessor {
  async publishOrdered(events: Event[]) {
    // Entity ID'ye göre sıralama için grupla
    const grouped = this.groupBy(events, e => e.entityId);

    for (const [entityId, entityEvents] of grouped) {
      // Timestamp'e göre sırala
      entityEvents.sort((a, b) => a.timestamp - b.timestamp);

      // Aynı partition key ile gönder
      for (const event of entityEvents) {
        await this.kafka.send({
          topic: 'events',
          key: entityId,  // Sıralamayı garanti eder
          value: event
        });
      }
    }
  }
}

Monitoring ve Gözlemlenebilirlik#

Takip Edilecek Temel Metrikler#

// Kapsamlı metrik toplama
class EventMetrics {
  private metrics = {
    messagesPublished: new Counter('messages_published_total'),
    messagesConsumed: new Counter('messages_consumed_total'),
    messagesFailed: new Counter('messages_failed_total'),
    processingDuration: new Histogram('message_processing_duration_seconds'),
    queueDepth: new Gauge('queue_depth'),
    consumerLag: new Gauge('consumer_lag'),
    dlqDepth: new Gauge('dlq_depth')
  };

  async recordProcessing(message: Message, processor: Function) {
    const timer = this.metrics.processingDuration.startTimer();

    try {
      const result = await processor(message);
      this.metrics.messagesConsumed.inc();
      return result;
    } catch (error) {
      this.metrics.messagesFailed.inc({
        error_type: error.constructor.name,
        queue: message.source
      });
      throw error;
    } finally {
      timer();
    }
  }
}

Sonuç#

Yönetilen varsayılan sistemlerin çoğunda geçerli: mesaj kalıbı, bulutunuzun zaten çalıştırdığı kuyruğu, topic’i ya da router’ı seçsin; kazandığınız eforu broker yerine teslimat semantiğine harcayın. Geçmiş event’leri replay etmeniz, partition’lanmış bir anahtar uzayında sıralama tutmanız ya da yönetilen kotayı aşan sürekli bir hacmi taşımanız gerekiyorsa varsayılanı geçin; Kafka ya da Kinesis operasyonel maliyetini tam da orada hak eder. Hangisinde karar kılarsanız kılın, DLQ’yu ve consumer lag alarmını ilk production mesajından önce kurun; hata senaryoları, ölçek senaryolarından çok daha erken gelir.


İlgili Derinlemesine İncelemeler:

Kaynaklar#

İlgili yazılar