İçeriğe atla

AWS Lambda, EventBridge ve DynamoDB ile CQRS

AWS Lambda, EventBridge ve DynamoDB ile pratik bir CQRS uygulaması; event sourcing, eventual consistency ve dağıtık sistemde debug yöntemleri.

Ayhan Sipahi Ayhan Sipahi

CQRS (Command Query Responsibility Segregation), yazma işlemlerini (commands) okuma işlemlerinden (queries) ayıran bir mimari pattern’dir. Aynı modeli iki iş için birden kullanmak yerine her taraf kendi işine uygun bir şekil alır: yazma tarafı doğrular ve saklar, okuma tarafı soruları hızlıca yanıtlar.

Ayrımı, iki iş yükünün şekli gerçekten farklıysa yap: 10:1 veya daha yüksek okuma-yazma oranı, farklı gecikme bütçeleri ya da yazma için optimize edilmiş bir şemanın karşılayamadığı query pattern’leri. Aşağıdaki kurulum her taraf için ayrı bir DynamoDB tablosu, her yol için ayrı bir Lambda fonksiyonu ve read model’i güncel tutan event’leri taşıyan EventBridge kullanıyor. Okuma ile yazma birbirine benziyorsa pattern yalnızca hareketli parça ekler.

Temel Prensip#

Geleneksel mimarilerde, genellikle hem okuma hem yazma için aynı data modelini kullanırsın:

// Geleneksel yaklaşım - her şey için aynı model
class OrderService {
  async createOrder(orderData) {
    // Aynı tabloya yaz
    return await db.orders.insert(orderData);
  }

  async getOrderHistory(customerId) {
    // Kompleks join'lerle aynı tablodan oku
    return await db.orders.find()
      .join('customers')
      .join('products')
      .where('customerId', customerId);
  }
}

CQRS ile bunu iki optimize edilmiş modele bölersin:

// CQRS yaklaşımı - ayrı optimize edilmiş modeller
class OrderCommandService {
  async createOrder(orderData) {
    // Write-optimized: Basit, hızlı insert'ler
    await writeDb.orders.insert(orderData);
    // Read model güncellemeleri için event publish et
    await eventBus.publish('OrderCreated', orderData);
  }
}

class OrderQueryService {
  async getOrderHistory(customerId) {
    // Read-optimized: Önceden hesaplanmış, denormalize edilmiş data
    return await readDb.customerOrderHistory.find(customerId);
  }
}

CQRS Neyi Çözüyor#

Takımları ayrıma iten dört somut problem var:

  1. Performans Uyumsuzluğu: Write’lar validation ve consistency, read’ler hız istiyor. Aynı model her ikisine hizmet edemez.
  2. Ölçek Uyumsuzluğu: Çoğu sistemde 10:1 veya 100:1 read-write oranı var; read tarafı çok daha fazla ölçeklenmeli.
  3. Model Karmaşıklığı: Write’lar için optimize etmek read’leri karmaşık yapıyor ve tersi. CQRS her tarafı ayrı optimize eder.
  4. Takım Paralelleştirmesi: Farklı takımlar read ve write taraflarında bağımsız çalışabilir; deploy’lar da birbirini bloke etmez.

CQRS Ne Zaman Mantıklı#

CQRS’i kullan eğer:

  • Yüksek read-write oranların var (10:1 veya daha fazla)
  • Read’ler ve write’lar için farklı performans gereksinimlerin var
  • Kompleks raporlama veya analitik ihtiyaçların var
  • Read’leri ve write’ları bağımsız scale etmen gerekiyor
  • Birden fazla data temsil ihtiyacın var (API’lar, raporlar, dashboard’lar)

CQRS’ten kaçın eğer:

  • Basit CRUD uygulamalarınız var
  • Düşük trafik uygulamalarınız var
  • Her yerde strong consistency gereksinimi var
  • Karmaşıklığı handle edemeyecek küçük takımınız var
  • Benzer read ve write pattern’leriniz var

CQRS’e Götüren Problem#

CQRS’e neden ihtiyaç duyulduğu somut bir örnekle açıklanabilir. Monolitik bir Lambda fonksiyonu her şeyi handle eder: ürün kataloğu okumaları, sipariş işleme, envanter güncellemeleri. Bir flash satış sırasında şu sorunlar ortaya çıkar:

  1. DynamoDB throttling: Siparişlerden gelen yoğun yazma işlemleri, gezinen kullanıcıların okuma işlemleriyle aynı kapasite için yarışır
  2. Lambda timeout’ları: Kompleks aggregation query’leri kayda değer süre alır
  3. Maliyet sorunu: Provisioned capacity yalnızca yoğun saatler için gerekir ama gün boyu ödenir
  4. Data tutarsızlığı: Concurrent update’ler envanter sayılarını bozar

Ürün detay sayfaları bu trafiğin büyük kısmını taşır ve yavaş kalır. Sebebi yapısaldır: sipariş işlemeye göre ayarlanmış bir data modelinden okurlar.

Mimarinin Evrimi#

CQRS Öncesi (Monolit):

// Her şeyi handle eden tek Lambda - ilk başta basit görünüyordu
export const handler = async (event: APIGatewayEvent) => {
  const { httpMethod, path } = event;

  if (httpMethod === 'GET' && path === '/products') {
    // 3 tabloyu join eden kompleks query
    const products = await dynamoClient.query({
      TableName: 'MainTable',
      IndexName: 'GSI1',
      KeyConditionExpression: 'GSI1PK = :pk',
      ExpressionAttributeValues: { ':pk': 'PRODUCT' }
    }).promise();

    // Sonra her ürün için envanter fetch et (N+1 query problemi)
    for (const product of products.Items) {
      const inventory = await getInventory(product.id);
      product.availableQuantity = inventory.quantity;
    }

    return { statusCode: 200, body: JSON.stringify(products) };
  }

  if (httpMethod === 'POST' && path === '/orders') {
    // Aynı tabloya yaz, throughput için yarış
    await createOrder(JSON.parse(event.body));
  }
};

CQRS Sonrası (Ayrılmış Concern’ler):

Command Side: Yazma Katmanı#

// commands/create-order.ts - Sadece sipariş işlemeye odaklanmış
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import { DynamoDBDocumentClient, PutCommand } from '@aws-sdk/lib-dynamodb';
import { EventBridgeClient, PutEventsCommand } from '@aws-sdk/client-eventbridge';
import { z } from 'zod';
import { ulid } from 'ulid';

const dynamoClient = DynamoDBDocumentClient.from(new DynamoDBClient({}));
const eventBridge = new EventBridgeClient({});

// Command sınırında Zod ile input validation
const CreateOrderSchema = z.object({
  customerId: z.string().uuid(),
  items: z.array(z.object({
    productId: z.string(),
    quantity: z.number().positive(),
    price: z.number().positive()
  })).min(1),
  shippingAddress: z.object({
    street: z.string(),
    city: z.string(),
    country: z.string(),
    postalCode: z.string()
  })
});

export const handler = async (event: any) => {
  // Input'u parse et ve validate et
  const input = CreateOrderSchema.parse(JSON.parse(event.body));

  const orderId = ulid(); // Time-sortable ID'ler log korelasyonunu kolaylaştırır
  const timestamp = Date.now();

  // Command store'a yaz (write-optimized tablo)
  const order = {
    PK: `ORDER#${orderId}`,
    SK: `ORDER#${orderId}`,
    id: orderId,
    customerId: input.customerId,
    items: input.items,
    total: input.items.reduce((sum, item) => sum + (item.price * item.quantity), 0),
    status: 'PENDING',
    createdAt: timestamp,
    updatedAt: timestamp,
    version: 1 // Optimistic locking race condition'ları önler
  };

  try {
    await dynamoClient.send(new PutCommand({
      TableName: process.env.WRITE_TABLE_NAME!,
      Item: order,
      ConditionExpression: 'attribute_not_exists(PK)' // Duplicate'leri önle
    }));

    // Read model update'leri için event publish et
    await eventBridge.send(new PutEventsCommand({
      Entries: [{
        Source: 'orders.service',
        DetailType: 'OrderCreated',
        Detail: JSON.stringify({
          orderId,
          customerId: input.customerId,
          items: input.items,
          total: order.total,
          timestamp
        }),
        EventBusName: process.env.EVENT_BUS_NAME
      }]
    }));

    return {
      statusCode: 201,
      body: JSON.stringify({ orderId, status: 'CREATED' })
    };

  } catch (error) {
    console.error('Sipariş oluşturma başarısız:', error);
    // Proper error handling ve compensation implement et
    throw error;
  }
};

Query Side: Optimize Edilmiş Okumalar#

// queries/get-product-catalog.ts - Performans için read-optimized
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import { DynamoDBDocumentClient, QueryCommand } from '@aws-sdk/lib-dynamodb';

const dynamoClient = DynamoDBDocumentClient.from(new DynamoDBClient({}));

// Hızlı okumalar için pre-computed, denormalized data
export const handler = async (event: any) => {
  const { category, limit = 20, lastKey } = event.queryStringParameters || {};

  // Pre-computed aggregation'larla read-optimized tablodan oku
  const response = await dynamoClient.send(new QueryCommand({
    TableName: process.env.READ_TABLE_NAME!,
    IndexName: 'CategoryIndex',
    KeyConditionExpression: 'category = :category',
    ExpressionAttributeValues: {
      ':category': category || 'ALL'
    },
    Limit: Number(limit),
    ExclusiveStartKey: lastKey ? JSON.parse(Buffer.from(lastKey, 'base64').toString()) : undefined,
    // Listing için sadece ihtiyacımız olanı fetch et
    ProjectionExpression: 'id, #n, price, imageUrl, averageRating, reviewCount, inStock',
    ExpressionAttributeNames: {
      '#n': 'name' // 'name' DynamoDB'de reserved word
    }
  }));

  return {
    statusCode: 200,
    headers: {
      'Cache-Control': 'public, max-age=300', // Ürün listelemeleri için 5 dakika cache
    },
    body: JSON.stringify({
      products: response.Items,
      nextKey: response.LastEvaluatedKey
        ? Buffer.from(JSON.stringify(response.LastEvaluatedKey)).toString('base64')
        : null
    })
  };
};

Event Processor: Model Senkronizasyonu#

Çoğu CQRS implementasyonu senkronizasyon katmanında başarısız olur:

// processors/sync-read-models.ts - Kritik senkronizasyon katmanı
import { EventBridgeEvent } from 'aws-lambda';
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import { DynamoDBDocumentClient, UpdateCommand } from '@aws-sdk/lib-dynamodb';
import { SQSClient, SendMessageCommand } from '@aws-sdk/client-sqs';

const dynamoClient = DynamoDBDocumentClient.from(new DynamoDBClient({}));
const sqsClient = new SQSClient({});

interface OrderCreatedEvent {
  orderId: string;
  customerId: string;
  items: Array<{ productId: string; quantity: number; price: number }>;
  total: number;
  timestamp: number;
}

export const handler = async (event: EventBridgeEvent<'OrderCreated', OrderCreatedEvent>) => {
  const { detail } = event;

  // Birden fazla read model'i paralel olarak güncelle
  const updatePromises = [];

  // 1. Müşteri sipariş geçmişini güncelle (müşteri query'leri için optimize edilmiş)
  updatePromises.push(
    dynamoClient.send(new UpdateCommand({
      TableName: process.env.READ_TABLE_NAME!,
      Key: {
        PK: `CUSTOMER#${detail.customerId}`,
        SK: `ORDER#${detail.timestamp}#${detail.orderId}`
      },
      UpdateExpression: 'SET orderId = :orderId, total = :total, #items = :items, createdAt = :timestamp',
      ExpressionAttributeNames: {
        '#items': 'items'
      },
      ExpressionAttributeValues: {
        ':orderId': detail.orderId,
        ':total': detail.total,
        ':items': detail.items,
        ':timestamp': detail.timestamp
      }
    }))
  );

  // 2. Ürün istatistiklerini güncelle (popüler ürünler, çok satanlar için)
  for (const item of detail.items) {
    updatePromises.push(
      dynamoClient.send(new UpdateCommand({
        TableName: process.env.READ_TABLE_NAME!,
        Key: {
          PK: `PRODUCT#${item.productId}`,
          SK: 'STATS'
        },
        UpdateExpression: `
          ADD salesCount :quantity, revenue :revenue
          SET lastSoldAt = :timestamp
        `,
        ExpressionAttributeValues: {
          ':quantity': item.quantity,
          ':revenue': item.price * item.quantity,
          ':timestamp': detail.timestamp
        }
      }))
    );
  }

  // 3. Günlük satış aggregation'larını güncelle (dashboard'lar için)
  const dateKey = new Date(detail.timestamp).toISOString().split('T')[0];
  updatePromises.push(
    dynamoClient.send(new UpdateCommand({
      TableName: process.env.READ_TABLE_NAME!,
      Key: {
        PK: `SALES#${dateKey}`,
        SK: 'AGGREGATE'
      },
      UpdateExpression: 'ADD orderCount :one, totalRevenue :total',
      ExpressionAttributeValues: {
        ':one': 1,
        ':total': detail.total
      }
    }))
  );

  try {
    await Promise.all(updatePromises);
  } catch (error) {
    console.error('Read modelleri güncellenemedi:', error);

    // Manuel müdahale için DLQ'ya gönder
    await sqsClient.send(new SendMessageCommand({
      QueueUrl: process.env.DLQ_URL!,
      MessageBody: JSON.stringify({
        event: 'OrderCreated',
        detail,
        error: error.message,
        timestamp: Date.now()
      })
    }));

    throw error; // Lambda'nın retry yapmasına izin ver
  }
};

CDK ile Infrastructure as Code#

İşte komple serverless CQRS setup’ı:

// infrastructure/cqrs-stack.ts
import { Stack, StackProps, Duration, RemovalPolicy } from 'aws-cdk-lib';
import { Construct } from 'constructs';
import * as lambda from 'aws-cdk-lib/aws-lambda-nodejs';
import * as dynamodb from 'aws-cdk-lib/aws-dynamodb';
import * as events from 'aws-cdk-lib/aws-events';
import * as targets from 'aws-cdk-lib/aws-events-targets';
import * as apigateway from 'aws-cdk-lib/aws-apigateway';
import * as sqs from 'aws-cdk-lib/aws-sqs';
import { Runtime } from 'aws-cdk-lib/aws-lambda';

export class CQRSServerlessStack extends Stack {
  constructor(scope: Construct, id: string, props?: StackProps) {
    super(scope, id, props);

    // Write model tablosu - write'lar için optimize edilmiş
    const writeTable = new dynamodb.Table(this, 'WriteTable', {
      partitionKey: { name: 'PK', type: dynamodb.AttributeType.STRING },
      sortKey: { name: 'SK', type: dynamodb.AttributeType.STRING },
      billingMode: dynamodb.BillingMode.PAY_PER_REQUEST, // On-demand: spike'larda throttling yok
      stream: dynamodb.StreamViewType.NEW_AND_OLD_IMAGES, // Change data capture için
      pointInTimeRecovery: true,
      removalPolicy: RemovalPolicy.RETAIN
    });

    // Read model tablosu - query'ler için optimize edilmiş
    const readTable = new dynamodb.Table(this, 'ReadTable', {
      partitionKey: { name: 'PK', type: dynamodb.AttributeType.STRING },
      sortKey: { name: 'SK', type: dynamodb.AttributeType.STRING },
      billingMode: dynamodb.BillingMode.PAY_PER_REQUEST,
      pointInTimeRecovery: true,
      removalPolicy: RemovalPolicy.RETAIN
    });

    // Farklı query pattern'leri için GSI'lar ekle
    readTable.addGlobalSecondaryIndex({
      indexName: 'CategoryIndex',
      partitionKey: { name: 'category', type: dynamodb.AttributeType.STRING },
      sortKey: { name: 'popularity', type: dynamodb.AttributeType.NUMBER },
      projectionType: dynamodb.ProjectionType.ALL
    });

    readTable.addGlobalSecondaryIndex({
      indexName: 'CustomerIndex',
      partitionKey: { name: 'customerId', type: dynamodb.AttributeType.STRING },
      sortKey: { name: 'createdAt', type: dynamodb.AttributeType.NUMBER },
      projectionType: dynamodb.ProjectionType.ALL
    });

    // CQRS event'leri için event bus
    const eventBus = new events.EventBus(this, 'CQRSEventBus', {
      eventBusName: 'cqrs-events'
    });

    // Başarısız event'ler için dead letter queue
    const dlq = new sqs.Queue(this, 'EventDLQ', {
      queueName: 'cqrs-event-dlq',
      retentionPeriod: Duration.days(14)
    });

    // Command handler'lar
    const createOrderHandler = new lambda.NodejsFunction(this, 'CreateOrderHandler', {
      entry: 'src/commands/create-order.ts',
      runtime: Runtime.NODEJS_22_X,
      memorySize: 1024,
      timeout: Duration.seconds(10),
      environment: {
        WRITE_TABLE_NAME: writeTable.tableName,
        EVENT_BUS_NAME: eventBus.eventBusName,
        AWS_NODEJS_CONNECTION_REUSE_ENABLED: '1'
      },
      bundling: {
        minify: true,
        target: 'node22',
        externalModules: ['@aws-sdk/*']
      }
    });

    writeTable.grantWriteData(createOrderHandler);
    eventBus.grantPutEventsTo(createOrderHandler);

    // Query handler'lar
    const getProductsHandler = new lambda.NodejsFunction(this, 'GetProductsHandler', {
      entry: 'src/queries/get-product-catalog.ts',
      runtime: Runtime.NODEJS_22_X,
      memorySize: 512, // Read-only, daha az memory gerekiyor
      timeout: Duration.seconds(5),
      environment: {
        READ_TABLE_NAME: readTable.tableName,
        AWS_NODEJS_CONNECTION_REUSE_ENABLED: '1'
      }
    });

    readTable.grantReadData(getProductsHandler);

    // Read model'leri senkronize etmek için event processor
    const syncProcessor = new lambda.NodejsFunction(this, 'SyncProcessor', {
      entry: 'src/processors/sync-read-models.ts',
      runtime: Runtime.NODEJS_22_X,
      memorySize: 2048, // Batch update'leri handle ediyor
      timeout: Duration.seconds(30),
      reservedConcurrentExecutions: 10, // Downstream servisleri overwhelm etmeyi önle
      environment: {
        READ_TABLE_NAME: readTable.tableName,
        DLQ_URL: dlq.queueUrl,
        AWS_NODEJS_CONNECTION_REUSE_ENABLED: '1'
      },
      deadLetterQueue: dlq,
      retryAttempts: 2
    });

    readTable.grantWriteData(syncProcessor);
    dlq.grantSendMessages(syncProcessor);

    // Event rule'ları
    new events.Rule(this, 'OrderCreatedRule', {
      eventBus,
      eventPattern: {
        source: ['orders.service'],
        detailType: ['OrderCreated']
      },
      targets: [new targets.LambdaFunction(syncProcessor, {
        retryAttempts: 2,
        maxEventAge: Duration.hours(2)
      })]
    });

    // API Gateway
    const api = new apigateway.RestApi(this, 'CQRSAPI', {
      restApiName: 'cqrs-api',
      defaultCorsPreflightOptions: {
        allowOrigins: apigateway.Cors.ALL_ORIGINS,
        allowMethods: apigateway.Cors.ALL_METHODS
      }
    });

    // Command endpoint'leri
    const orders = api.root.addResource('orders');
    orders.addMethod('POST', new apigateway.LambdaIntegration(createOrderHandler));

    // Query endpoint'leri
    const products = api.root.addResource('products');
    products.addMethod('GET', new apigateway.LambdaIntegration(getProductsHandler));
  }
}

Eventual Consistency Yönetimi#

CQRS, eventual consistency’yi kabul etmek demek. Kullanıcının kafasını karıştırmadan bunu yönetmenin üç yolu:

// strategies/consistency-handling.ts
export class ConsistencyStrategy {
  // Strateji 1: Optimistic UI update'leri
  async createOrderWithOptimisticUpdate(orderData: any) {
    // Kullanıcıya hemen başarı göster
    const tempOrderId = `temp_${Date.now()}`;
    updateUI({ orderId: tempOrderId, status: 'processing' });

    try {
      const response = await fetch('/api/orders', {
        method: 'POST',
        body: JSON.stringify(orderData)
      });

      const { orderId } = await response.json();

      // Temp ID'yi gerçek ID ile değiştir
      updateUI({ oldId: tempOrderId, newId: orderId, status: 'confirmed' });

      // Read model update'i için poll et
      await this.waitForReadModelSync(orderId);

    } catch (error) {
      // Optimistic update'i geri al
      removeFromUI(tempOrderId);
      showError('Sipariş başarısız');
    }
  }

  // Strateji 2: Exponential backoff ile polling
  async waitForReadModelSync(orderId: string, maxAttempts = 5) {
    let attempts = 0;
    let delay = 100; // 100ms ile başla

    while (attempts < maxAttempts) {
      const order = await this.checkReadModel(orderId);

      if (order) {
        return order;
      }

      await new Promise(resolve => setTimeout(resolve, delay));
      delay *= 2; // Exponential backoff
      attempts++;
    }

    // Command model query'sine fall back et
    return this.queryCommandModel(orderId);
  }

  // Strateji 3: WebSocket notification'ları
  subscribeToOrderUpdates(customerId: string) {
    const ws = new WebSocket(`wss://api.example.com/orders/${customerId}`);

    ws.onmessage = (event) => {
      const update = JSON.parse(event.data);
      if (update.type === 'READ_MODEL_SYNCED') {
        refreshOrderList();
      }
    };
  }
}

Serverless’ta CQRS Test Etme#

Distributed sistemleri test etmek zor. İşte pratik bir yaklaşım:

// tests/cqrs-integration.test.ts
import { EventBridgeClient, PutEventsCommand } from '@aws-sdk/client-eventbridge';
import { DynamoDBDocumentClient, UpdateCommand } from '@aws-sdk/lib-dynamodb';
import { SQSClient } from '@aws-sdk/client-sqs';
import { mockClient } from 'aws-sdk-client-mock';
import { handler } from '../src/commands/create-order';
import { handler as syncProcessor } from '../src/processors/sync-read-models';

describe('CQRS Event Flow', () => {
  const eventBridgeMock = mockClient(EventBridgeClient);
  const dynamoMock = mockClient(DynamoDBDocumentClient);
  const sqsMock = mockClient(SQSClient);

  beforeEach(() => {
    eventBridgeMock.reset();
    dynamoMock.reset();
    sqsMock.reset();
  });

  test('Sipariş oluşturma read model güncellemesini tetikler', async () => {
    // Arrange
    const orderId = 'test-order-123';
    eventBridgeMock.on(PutEventsCommand).resolves({
      FailedEntryCount: 0,
      Entries: [{ EventId: 'event-123' }]
    });

    // Act - Sipariş oluştur
    const response = await handler({
      body: JSON.stringify({
        customerId: 'customer-123',
        items: [{ productId: 'prod-1', quantity: 2, price: 99.99 }]
      })
    });

    // Assert - Event publish edildi
    expect(eventBridgeMock.calls()).toHaveLength(1);
    const eventCall = eventBridgeMock.call(0);
    expect(eventCall.args[0].input.Entries[0].DetailType).toBe('OrderCreated');

    // Event processor'ı simüle et
    await syncProcessor({
      detail: JSON.parse(eventCall.args[0].input.Entries[0].Detail)
    });

    // Assert - Read model'ler güncellendi
    const readModelCalls = dynamoMock.calls().filter(
      call => call.args[0].input.TableName === 'ReadTable'
    );
    expect(readModelCalls).toHaveLength(3); // Customer, Product, Daily stats
  });

  test('Başarısız event processing DLQ mesajı üretir', async () => {
    // DynamoDB failure'ı simüle et
    dynamoMock.on(UpdateCommand).rejects(new Error('Throttled'));

    const event = {
      detail: {
        orderId: 'order-123',
        customerId: 'customer-123',
        items: [],
        total: 100,
        timestamp: Date.now()
      }
    };

    await expect(syncProcessor(event)).rejects.toThrow('Throttled');

    // DLQ mesajını verify et
    const sqsCalls = sqsMock.calls();
    expect(sqsCalls).toHaveLength(1);
    expect(JSON.parse(sqsCalls[0].args[0].input.MessageBody))
      .toHaveProperty('error', 'Throttled');
  });
});

CQRS Monitoring ve Debugging#

CQRS’in distributed doğası debugging’i zorlaştırır. İşte etkili bir monitoring setup’ı:

// monitoring/cqrs-metrics.ts
import { MetricUnit, Metrics } from '@aws-lambda-powertools/metrics';
import { Tracer } from '@aws-lambda-powertools/tracer';
import { Logger } from '@aws-lambda-powertools/logger';

const metrics = new Metrics({ namespace: 'CQRS', serviceName: 'orders' });
const tracer = new Tracer({ serviceName: 'orders' });
const logger = new Logger({ serviceName: 'orders' });

export const instrumentedHandler = tracer.captureLambdaHandler(
  metrics.logMetrics(
    async (event: any) => {
      const segment = tracer.getSegment();

      // Command/query separation'ı track et
      const operationType = event.httpMethod === 'GET' ? 'QUERY' : 'COMMAND';
      metrics.addMetric(`${operationType}_REQUEST`, MetricUnit.Count, 1);

      const startTime = Date.now();

      try {
        // Servisler arası tracing için correlation ID ekle
        const correlationId = event.headers['x-correlation-id'] || ulid();
        segment?.addAnnotation('correlationId', correlationId);
        logger.appendKeys({ correlationId });

        // Read/write model sync lag'i track et
        if (operationType === 'QUERY') {
          const syncLag = await measureSyncLag();
          metrics.addMetric('READ_MODEL_LAG_MS', MetricUnit.Milliseconds, syncLag);

          if (syncLag > 5000) {
            logger.warn('Yüksek read model lag tespit edildi', { syncLag });
          }
        }

        const result = await processRequest(event);

        metrics.addMetric(`${operationType}_SUCCESS`, MetricUnit.Count, 1);
        metrics.addMetric(`${operationType}_DURATION`, MetricUnit.Milliseconds,
          Date.now() - startTime);

        return result;

      } catch (error) {
        metrics.addMetric(`${operationType}_ERROR`, MetricUnit.Count, 1);
        logger.error('Request başarısız', { error, event });
        throw error;
      }
    }
  )
);

// Custom CloudWatch dashboard
export const dashboardConfig = {
  widgets: [
    {
      type: 'metric',
      properties: {
        metrics: [
          ['CQRS', 'COMMAND_REQUEST', { stat: 'Sum' }],
          ['.', 'QUERY_REQUEST', { stat: 'Sum' }],
          ['.', 'READ_MODEL_LAG_MS', { stat: 'Average' }]
        ],
        period: 300,
        stat: 'Average',
        region: 'us-east-1',
        title: 'CQRS Operations'
      }
    }
  ]
};

Maliyet Nerede Değişiyor#

Modelleri ayırmak faturanın toplamından çok şeklini değiştirir. Farkın büyük kısmını dört etki açıklar:

  1. Peak için over-provisioning yok: iki tabloda da on-demand faturalama, provisioned kurulumun genelde en büyük kalemi olan peak kapasite primini ortadan kaldırır
  2. Cache’lenmiş read model’ler tekrar eden trafiği veritabanına ulaşmadan karşılar
  3. Odaklanmış fonksiyonlar, en ağır dalına göre boyutlanmış tek bir handler’dan daha az memory ister
  4. Purpose-built index’ler okuma yolundaki scan ve N+1 sorgularının yerini alır

Bunun karşılığında yazma yolu iki kez ödeme yapar: bir kez command tablosu, bir kez de read tablosuna yazılan projeksiyon için. Okuma tarafı baskınsa bu takas kazançlıdır, değilse değildir. Karar vermeden önce iki yolu da kendi trafik dağılımınıza göre fiyatlandırın.

Öğrenilen Dersler#

1. Tek Read Model ile Başla#

Erken implementasyonlar genelde fazla sayıda read model üretir. Tek model ile başla, yenisini ancak bir query pattern gerektirdiğinde ekle. Her ek model, doğru kalması gereken ek sync mantığı demektir.

2. Event Versioning Kritik#

Başlangıçta event’leri version’lamamak yaygın bir hatadır: OrderCreated’a bir field eklendiğinde her consumer bozulur. Doğru yaklaşım:

interface OrderCreatedV1 {
  version: 1;
  orderId: string;
  customerId: string;
  total: number;
}

interface OrderCreatedV2 {
  version: 2;
  orderId: string;
  customerId: string;
  total: number;
  currency: string; // Yeni field
}

// Handler her iki version'ı da destekliyor
export const handler = async (event: OrderCreatedV1 | OrderCreatedV2) => {
  const currency = 'version' in event && event.version >= 2
    ? (event as OrderCreatedV2).currency
    : 'USD'; // V1 için default
};

3. Her Yerde Idempotency#

Event’ler birden fazla kez deliver edilebilir. Her handler idempotent olmalı:

// Idempotency sağlamak için conditional write'lar kullan
await dynamoClient.send(new PutCommand({
  TableName: TABLE_NAME,
  Item: processedEvent,
  ConditionExpression: 'attribute_not_exists(eventId)'
}));

4. Sync Lag’i Monitor Et#

Command execution ve read model update arasındaki süre en önemli metriktir. 5 saniyeyi aşarsa alert verilmelidir.

5. Reconciliation İçin Plan Yap#

Read model’ler drift edecek. Command ve query model’leri karşılaştıran, tutarsızlıkları düzelten bir nightly job bu sorunu çözer:

// Düşük trafikli saatlerde günde bir çalışır
export const reconciliationJob = async () => {
  const commandRecords = await scanCommandTable();
  const readRecords = await scanReadTable();

  const discrepancies = findDiscrepancies(commandRecords, readRecords);

  for (const issue of discrepancies) {
    await republishEvent(issue.originalEvent);
    logger.warn('Reconciliation gerekli', { issue });
  }

  metrics.addMetric('RECONCILIATION_FIXES', MetricUnit.Count, discrepancies.length);
};

CQRS’ten Kaçınma Durumları#

CQRS karmaşıklık ekler. Doğru bağlamda değer; ama şu durumlarda kaçının:

  1. Read/write pattern’leriniz benzer
  2. Basit CRUD işlemleriniz var
  3. Her yerde strong consistency gerekli
  4. Takımınız eventual consistency ile rahat değil
  5. Performance sorunları yaşamıyorsunuz

Düz CRUD ekranlarından oluşan bir iç admin panelinde CQRS başarısızlık modudur: karşılığı olmayan karmaşıklık.

Sessizce Duran Event Processor#

En çok canını yakan arıza, sessizce duran event processor’dur. Command tarafı yazmayı kabul etmeye devam eder, read model ilerlemez ve müşteriler kimse fark etmeden eski sipariş durumlarını görür. Yavaş bir event kategorisinde alınan Lambda timeout’u ile kimsenin bakmadığı bir dead letter queue birleşince tam olarak bu tablo çıkar.

Bunu önleyen dört yapılandırma detayı var:

  1. DLQ’yu mesaj gelişinde tetiklenen bir CloudWatch alarmıyla ve retry penceresinden uzun bir visibility timeout ile kur
  2. Event processor’a exponential backoff ve circuit breaker ekle
  3. Gecikmeyi kaldıramayan, kullanıcıya dönük işlemlerde read-after-write consistency kontrolü çalıştır
  4. Read model geride kaldığında “command model’e düş” yolunu hazır tut

Varsayılanın Sınırları#

Ayrı tablolar, ayrı fonksiyonlar, aradaki EventBridge: bu şekil, okuma tarafı yazma tarafından daha yoğun olduğunda ve aynı veriye farklı sorular sorduğunda geçerli. Tek modelle başlayın, nerede zorlandığını ölçün ve sadece zorlanan yolu ayırın. Başlangıç için tek read model yeter; ikincisi diyagram simetrik dursun diye değil, bir query gerektirdiği için gelmeli.

Strong consistency bir tercih değil de ürün gereksinimiyse varsayılanı geçersiz kılın. Kendi yazdığını anında görmesi gereken kullanıcı, read model’in biraz sonra yetişmesiyle ilgilenmez. O yolu command model’den servis edin, geri kalanı CQRS’te bırakın; eventual consistency kullanıcıya görünen bir problem olmaktan çıkar.

Kaynaklar#

İlgili yazılar