Saga Pattern ile Dağıtık Transaction'lar: ACID Olmadan Consistency Sağlamak
Microservices mimarisinde AWS Step Functions ve EventBridge kullanarak Saga pattern implementasyonu: idempotency, compensation logic ve production-ready pattern'ler.
Microservices tek bir iş işlemini birden fazla database’e böler ve ACID garantileri servis sınırında biter. Saga pattern bu boşluğu local transaction’lar ve compensating transaction’larla dolduruyor: commit edilmiş her adım, sonraki bir adım başarısız olduğunda geri alınabiliyor ve sistem eventual consistency’ye oturuyor.
Sipariş benzeri akışlar için güvenli varsayılan AWS Step Functions ile orchestration: saga state’i ve hata yolları tek yerde görünür kalıyor. EventBridge choreography daha ucuz ve daha gevşek bağlı; üç-dört otonom servis zaten domain event paylaşıyorsa yerini hak ediyor. Hangi koordinasyon stilini seçersen seç, pattern’i retry’lar altında ayakta tutan şey idempotent step’ler, ters sırada compensation ve semantic locking.
Dağıtık Transaction Problemi#
Geleneksel ACID transaction’lar servis sınırları arasında çalışmaz; two-phase commit (2PC) ise atomicity’yi sıkı coupling ve tek bir failure noktası karşılığında satın alır.
Tipik bir e-commerce sipariş akışını düşün: sipariş oluşturma, stok rezervasyonu, ödeme işleme ve kargo onayı. Her adım farklı bir microservice ve kendi database’ine sahip. Stok rezerve edildikten sonra ödeme işleme başarısız olursa, stok rezervasyonunu geri almanın güvenilir bir yoluna ihtiyacın var. Doğru pattern’ler olmadan servisler birbirinden kopar: siparişler oluşur, ödemeler başarısız olur, stok geri verilmez. Ağ timeout’ları, provider rate limit’leri ve geçici servis kesintileri hep bu şekilde partial failure üretir; bu yüzden compensation logic akışın ilk sürümüyle birlikte gelmeli.
Compensation storage seviyesinde değil, iş seviyesinde çalışır. Forward step zaten commit edilmiştir ve etkisi başka okuyuculara görünür olabilir; bu yüzden rezervasyonu serbest bırakmak, ödemeyi iade etmek ve siparişi iptal olarak işaretlemek explicit birer compensating transaction olur. Bir saga’da işin büyük kısmı bu geri alma adımlarını tasarlamaktır.
Saga Pattern Temelleri#
Bir saga, şu şekilde işleyen bir dizi local transaction’dır:
- Her transaction tek bir servis içindeki data’yı günceller
- Her transaction bir sonraki transaction’ı tetiklemek için event veya message publish eder
- Bir transaction başarısız olursa, saga tamamlanan adımları geri almak için compensating transaction’ları çalıştırır
- Sistem anlık ACID consistency yerine eventual consistency sağlar
Saga’ları çalıştıran temel özellikler:
- Eventual Consistency: Sistem, tüm adımlar oturduktan sonra consistent duruma yakınsar
- Local Transaction’lar: Her servis kendi database’ini local ACID garantileriyle yönetir
- Compensating Transaction’lar: Her forward step için explicit rollback logic
- Idempotency: Tüm saga step’leri retry-safe olmalı
- Observability: Distributed flow’ları debug etmek için kritik; correlation ID ve structured logging her step’te bulunmalı
Saga pattern’i ne zaman kullanmalısın:
- Birden fazla database’i olan microservices mimarisi
- Birden fazla servisi kapsayan business process’ler
- Performance veya coupling endişeleri nedeniyle distributed transaction’lar kullanılamıyor
- Geçici inconsistency kabul edilebilir
- Tipik olarak maksimum 3-5 step (bunun ötesinde complexity hızla artar)
Orchestration vs Choreography#
Saga’ları koordine etmenin iki ana yaklaşımı var: orchestration ve choreography. Hangisini ne zaman kullanacağını anlamak başarılı implementasyon için kritik.
Choreography: Event-Driven Koordinasyon#
Choreography’de servisler merkezi bir coordinator olmadan domain event’ler aracılığıyla koordine olur. Her servis bir event aldığında ne yapacağını bilir.
Avantajları:
- Servisler arasında loose coupling
- Single point of failure yok
- Independent servisler için iyi scale eder
- Event-driven architecture için doğal fit
Dezavantajları:
- Control flow tek bir yerde görünmüyor
- Tüm saga flow’unu anlamak zor
- Debugging complexity (distributed logic)
- Saga state’i track etmek daha zor
- Cyclic dependency riski
En uygun durumlar:
- Maksimum 3-4 servis
- Independent, autonomous servisler
- Mevcut event-driven architecture
- Basit linear flow’lar
Orchestration: Merkezi Koordinasyon#
Orchestration’da merkezi bir coordinator (tipik olarak AWS Step Functions) saga flow’unu yönetir, her servise ne yapacağını söyler.
Avantajları:
- Net control flow visualization
- Daha kolay debugging ve monitoring
- Centralized error handling
- Built-in state management
- Complex flow’lar için daha iyi
Dezavantajları:
- Orchestrator coordination point
- Servisler orchestrator’a coupled
- Orchestrator tüm servisleri bilmeli
En uygun durumlar:
- Kompleks multi-step workflow’lar
- Saga state’e visibility gerekli
- Human approval step’leri var
- 4’ten fazla servis involved
- Strict ordering gereksinimleri
Karar Ağacı#
Yaklaşımı seçerken şu sırayla ilerle:
AWS Step Functions ile Orchestration Implementasyonu#
Aşağıda AWS CDK kullanarak production-ready bir e-commerce sipariş saga implementasyonu ele alınıyor.
Altyapı Kurulumu#
import * as cdk from 'aws-cdk-lib';
import * as sfn from 'aws-cdk-lib/aws-stepfunctions';
import * as tasks from 'aws-cdk-lib/aws-stepfunctions-tasks';
import * as lambda from 'aws-cdk-lib/aws-lambda';
import * as dynamodb from 'aws-cdk-lib/aws-dynamodb';
import { Construct } from 'constructs';
export class OrderSagaStack extends cdk.Stack {
constructor(scope: Construct, id: string) {
super(scope, id);
// Saga state persistence için DynamoDB table'lar
const ordersTable = new dynamodb.Table(this, 'Orders', {
partitionKey: { name: 'orderId', type: dynamodb.AttributeType.STRING },
sortKey: { name: 'transactionId', type: dynamodb.AttributeType.STRING },
billingMode: dynamodb.BillingMode.PAY_PER_REQUEST
});
const inventoryTable = new dynamodb.Table(this, 'Inventory', {
partitionKey: { name: 'productId', type: dynamodb.AttributeType.STRING },
billingMode: dynamodb.BillingMode.PAY_PER_REQUEST
});
const paymentsTable = new dynamodb.Table(this, 'Payments', {
partitionKey: { name: 'paymentId', type: dynamodb.AttributeType.STRING },
billingMode: dynamodb.BillingMode.PAY_PER_REQUEST
});
// Forward transaction Lambda'lar
const createOrder = new lambda.Function(this, 'CreateOrder', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'createOrder.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
ORDERS_TABLE: ordersTable.tableName
}
});
const reserveInventory = new lambda.Function(this, 'ReserveInventory', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'reserveInventory.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
INVENTORY_TABLE: inventoryTable.tableName
}
});
const processPayment = new lambda.Function(this, 'ProcessPayment', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'processPayment.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
PAYMENTS_TABLE: paymentsTable.tableName
}
});
const confirmOrder = new lambda.Function(this, 'ConfirmOrder', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'confirmOrder.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
ORDERS_TABLE: ordersTable.tableName
}
});
// Compensating transaction Lambda'lar (rollback)
const cancelOrder = new lambda.Function(this, 'CancelOrder', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'cancelOrder.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
ORDERS_TABLE: ordersTable.tableName
}
});
const releaseInventory = new lambda.Function(this, 'ReleaseInventory', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'releaseInventory.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
INVENTORY_TABLE: inventoryTable.tableName
}
});
const refundPayment = new lambda.Function(this, 'RefundPayment', {
runtime: lambda.Runtime.NODEJS_20_X,
handler: 'refundPayment.handler',
code: lambda.Code.fromAsset('lambda'),
environment: {
PAYMENTS_TABLE: paymentsTable.tableName
}
});
// Table permission'ları
ordersTable.grantReadWriteData(createOrder);
ordersTable.grantReadWriteData(confirmOrder);
ordersTable.grantReadWriteData(cancelOrder);
inventoryTable.grantReadWriteData(reserveInventory);
inventoryTable.grantReadWriteData(releaseInventory);
paymentsTable.grantReadWriteData(processPayment);
paymentsTable.grantReadWriteData(refundPayment);
// Retry configuration ile Step Functions task'lar
const createOrderTask = new tasks.LambdaInvoke(this, 'CreateOrderTask', {
lambdaFunction: createOrder,
outputPath: '$.Payload',
retryOnServiceExceptions: true
}).addRetry({
errors: ['ThrottlingException', 'ServiceUnavailable'],
interval: cdk.Duration.seconds(2),
maxAttempts: 3,
backoffRate: 2.0
});
const reserveInventoryTask = new tasks.LambdaInvoke(this, 'ReserveInventoryTask', {
lambdaFunction: reserveInventory,
outputPath: '$.Payload',
retryOnServiceExceptions: true
}).addRetry({
errors: ['ThrottlingException'],
interval: cdk.Duration.seconds(1),
maxAttempts: 3,
backoffRate: 1.5
});
const processPaymentTask = new tasks.LambdaInvoke(this, 'ProcessPaymentTask', {
lambdaFunction: processPayment,
outputPath: '$.Payload',
retryOnServiceExceptions: true,
timeout: cdk.Duration.seconds(30)
}).addRetry({
errors: ['PaymentProcessingException'],
interval: cdk.Duration.seconds(3),
maxAttempts: 2,
backoffRate: 2.0
});
const confirmOrderTask = new tasks.LambdaInvoke(this, 'ConfirmOrderTask', {
lambdaFunction: confirmOrder,
outputPath: '$.Payload'
});
// Compensation task'lar
const cancelOrderTask = new tasks.LambdaInvoke(this, 'CancelOrderTask', {
lambdaFunction: cancelOrder,
outputPath: '$.Payload'
});
const releaseInventoryTask = new tasks.LambdaInvoke(this, 'ReleaseInventoryTask', {
lambdaFunction: releaseInventory,
outputPath: '$.Payload'
});
const refundPaymentTask = new tasks.LambdaInvoke(this, 'RefundPaymentTask', {
lambdaFunction: refundPayment,
outputPath: '$.Payload'
});
// Compensation logic ile saga orchestration
// Compensation chain'i backwards oluştur
const compensatePayment = refundPaymentTask
.next(releaseInventoryTask)
.next(cancelOrderTask)
.next(new sfn.Fail(this, 'OrderFailed', {
error: 'OrderProcessingFailed',
cause: 'Payment processing failed, all steps compensated'
}));
const compensateInventory = releaseInventoryTask
.next(cancelOrderTask)
.next(new sfn.Fail(this, 'InventoryFailed', {
error: 'InventoryReservationFailed',
cause: 'Inventory reservation failed, order cancelled'
}));
// Compensation için catch block'larla forward flow oluştur
const sagaDefinition = createOrderTask
.next(reserveInventoryTask
.addCatch(compensateInventory, {
errors: ['States.ALL'],
resultPath: '$.inventoryError'
})
)
.next(processPaymentTask
.addCatch(compensatePayment, {
errors: ['States.ALL'],
resultPath: '$.paymentError'
})
)
.next(confirmOrderTask)
.next(new sfn.Succeed(this, 'OrderSuccess'));
// Saga state machine oluştur
const orderSaga = new sfn.StateMachine(this, 'OrderSaga', {
definition: sagaDefinition,
stateMachineType: sfn.StateMachineType.STANDARD,
timeout: cdk.Duration.minutes(5),
tracingEnabled: true
});
new cdk.CfnOutput(this, 'OrderSagaArn', {
value: orderSaga.stateMachineArn,
description: 'Order Saga State Machine ARN'
});
}
}
Bu infrastructure, doğru compensation chain’leriyle complete bir order processing saga setup’ı yapıyor. Her step’in transient error’lar için retry configuration’ı olduğuna ve compensation flow’larının reverse order’da build edildiğine dikkat et; bu tamamlanmış step’leri doğru şekilde geri almak için kritik.
Idempotency Implementasyonu#
Saga’larda idempotency tartışmasız gerekli. Step’ler retry’lar, failure’lar veya network sorunları nedeniyle birden fazla kez execute olabilir. Aşağıda düzgün idempotent operation implementasyonu gösteriliyor:
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import { DynamoDBDocumentClient, UpdateCommand, GetCommand } from '@aws-sdk/lib-dynamodb';
const client = new DynamoDBClient({});
const docClient = DynamoDBDocumentClient.from(client);
interface ReserveInventoryEvent {
orderId: string;
productId: string;
quantity: number;
transactionId: string; // Idempotency key
}
export const handler = async (event: ReserveInventoryEvent) => {
console.log('Reserving inventory', { orderId: event.orderId, productId: event.productId });
try {
// Idempotent check: Bu transaction zaten tamamlandı mı?
const existingReservation = await docClient.send(
new GetCommand({
TableName: process.env.INVENTORY_TABLE!,
Key: { productId: event.productId }
})
);
const reservation = existingReservation.Item?.reservations?.[event.transactionId];
if (reservation?.status === 'RESERVED') {
console.log('Inventory already reserved for this transaction', {
transactionId: event.transactionId
});
return {
success: true,
message: 'Idempotent: Already reserved',
reservationId: event.transactionId
};
}
// Conditional expression ile atomic update (semantic lock)
const result = await docClient.send(
new UpdateCommand({
TableName: process.env.INVENTORY_TABLE!,
Key: { productId: event.productId },
UpdateExpression:
'SET availableQuantity = availableQuantity - :qty, ' +
'reservedQuantity = reservedQuantity + :qty, ' +
'reservations.#txId = :reservation',
ConditionExpression: 'availableQuantity >= :qty',
ExpressionAttributeNames: {
'#txId': event.transactionId
},
ExpressionAttributeValues: {
':qty': event.quantity,
':reservation': {
orderId: event.orderId,
quantity: event.quantity,
status: 'RESERVED',
timestamp: Date.now(),
expiresAt: Date.now() + (15 * 60 * 1000) // 15 dakika timeout
}
},
ReturnValues: 'ALL_NEW'
})
);
console.log('Inventory reserved successfully', {
orderId: event.orderId,
reservationId: event.transactionId
});
return {
success: true,
reservationId: event.transactionId,
availableQuantity: result.Attributes?.availableQuantity
};
} catch (error: any) {
if (error.name === 'ConditionalCheckFailedException') {
console.error('Insufficient inventory', {
productId: event.productId,
requested: event.quantity
});
throw new Error('InsufficientInventory');
}
console.error('Failed to reserve inventory', error);
throw error;
}
};
Başlangıçtaki idempotency check, bu function aynı transactionId ile birden fazla kez execute olursa side effect olmadan aynı sonucu döndürmesini sağlıyor. Conditional expression atomic bir “semantic lock” sağlıyor; sadece yeterli quantity varsa stok rezerve ediyor.
Compensation: Stok Serbest Bırakma#
export const releaseInventoryHandler = async (event: ReserveInventoryEvent) => {
console.log('Releasing inventory reservation', {
orderId: event.orderId,
transactionId: event.transactionId
});
try {
const reservation = await docClient.send(
new GetCommand({
TableName: process.env.INVENTORY_TABLE!,
Key: { productId: event.productId }
})
);
const existingReservation = reservation.Item?.reservations?.[event.transactionId];
if (!existingReservation) {
console.log('Idempotent: Reservation already released or never existed', {
transactionId: event.transactionId
});
return { success: true, message: 'Already released' };
}
// Conditional check ile atomic release
await docClient.send(
new UpdateCommand({
TableName: process.env.INVENTORY_TABLE!,
Key: { productId: event.productId },
UpdateExpression:
'SET availableQuantity = availableQuantity + :qty, ' +
'reservedQuantity = reservedQuantity - :qty ' +
'REMOVE reservations.#txId',
ConditionExpression: 'attribute_exists(reservations.#txId)',
ExpressionAttributeNames: {
'#txId': event.transactionId
},
ExpressionAttributeValues: {
':qty': existingReservation.quantity
}
})
);
console.log('Inventory released successfully', { transactionId: event.transactionId });
return { success: true };
} catch (error: any) {
if (error.name === 'ConditionalCheckFailedException') {
// Reservation zaten released - idempotent success
return { success: true, message: 'Already released' };
}
throw error;
}
};
Compensation’ın da idempotent olduğuna dikkat et. Reservation yoksa success dönüyoruz; istenen end state zaten sağlanmış.
Idempotency ile Ödeme İşleme#
interface ProcessPaymentEvent {
orderId: string;
customerId: string;
amount: number;
currency: string;
paymentMethodId: string;
idempotencyKey: string;
}
export const processPaymentHandler = async (event: ProcessPaymentEvent) => {
console.log('Payment işleniyor', {
orderId: event.orderId,
amount: event.amount
});
try {
// Idempotency kontrolü: Bu ödeme daha önce işlendi mi?
const existingPayment = await docClient.send(
new GetCommand({
TableName: process.env.PAYMENTS_TABLE!,
Key: { paymentId: event.idempotencyKey }
})
);
if (existingPayment.Item?.status === 'COMPLETED') {
console.log('Idempotent: Ödeme zaten işlendi', {
paymentId: event.idempotencyKey
});
return {
success: true,
paymentId: event.idempotencyKey,
transactionId: existingPayment.Item.transactionId,
message: 'Already processed'
};
}
// Önce PENDING status ile ödeme kaydı oluştur
const paymentId = event.idempotencyKey;
await docClient.send(
new UpdateCommand({
TableName: process.env.PAYMENTS_TABLE!,
Key: { paymentId },
UpdateExpression:
'SET #status = :pending, orderId = :orderId, ' +
'amount = :amount, createdAt = :timestamp',
ConditionExpression: 'attribute_not_exists(paymentId)',
ExpressionAttributeNames: { '#status': 'status' },
ExpressionAttributeValues: {
':pending': 'PENDING',
':orderId': event.orderId,
':amount': event.amount,
':timestamp': Date.now()
}
})
);
// Harici payment provider'ı çağır (Stripe vb.)
const paymentResult = await callPaymentProvider({
amount: event.amount,
currency: event.currency,
customerId: event.customerId,
paymentMethodId: event.paymentMethodId,
idempotencyKey: event.idempotencyKey
});
if (!paymentResult.success) {
await docClient.send(
new UpdateCommand({
TableName: process.env.PAYMENTS_TABLE!,
Key: { paymentId },
UpdateExpression:
'SET #status = :failed, failureReason = :reason, updatedAt = :timestamp',
ExpressionAttributeNames: { '#status': 'status' },
ExpressionAttributeValues: {
':failed': 'FAILED',
':reason': paymentResult.error,
':timestamp': Date.now()
}
})
);
throw new Error(`PaymentFailed: ${paymentResult.error}`);
}
// COMPLETED status'a güncelle
await docClient.send(
new UpdateCommand({
TableName: process.env.PAYMENTS_TABLE!,
Key: { paymentId },
UpdateExpression:
'SET #status = :completed, transactionId = :txId, updatedAt = :timestamp',
ExpressionAttributeNames: { '#status': 'status' },
ExpressionAttributeValues: {
':completed': 'COMPLETED',
':txId': paymentResult.transactionId,
':timestamp': Date.now()
}
})
);
return {
success: true,
paymentId,
transactionId: paymentResult.transactionId
};
} catch (error: any) {
if (error.name === 'ConditionalCheckFailedException') {
// Başka bir execution zaten ödeme kaydını oluşturmuş
throw new Error('ConcurrentPaymentAttempt');
}
throw error;
}
};
// Simüle edilmiş payment provider çağrısı
async function callPaymentProvider(params: any) {
// Production'da: Stripe, Adyen vb. kendi idempotency key'leriyle
return {
success: Math.random() > 0.1, // %90 başarı oranı
transactionId: `txn_${Date.now()}`,
error: 'CardDeclined'
};
}
Bu ödeme implementasyonu üç aşamalı idempotency gösteriyor: mevcut completion kontrolü, PENDING kaydı oluşturma, ardından final state’e güncelleme. Bu pattern, fonksiyon birden fazla kez çalışsa bile müşteriden çift ücret alınmamasını sağlıyor.
EventBridge ile Choreography Implementasyonu#
Daha basit flow’lar için choreography daha iyi decoupling sağlayabilir. Aşağıda event-driven saga koordinasyonu gösteriliyor:
import { EventBridgeClient, PutEventsCommand } from '@aws-sdk/client-eventbridge';
// Order Service: Sipariş oluşturur ve OrderCreated event'i publish eder
export const createOrderHandler = async (event: any) => {
const orderId = `ord_${Date.now()}`;
// Database'de sipariş oluştur
await saveOrder({
orderId,
customerId: event.customerId,
items: event.items,
status: 'PENDING'
});
// OrderCreated event'i publish et
const eventBridge = new EventBridgeClient({});
await eventBridge.send(
new PutEventsCommand({
Entries: [{
Source: 'order.service',
DetailType: 'OrderCreated',
Detail: JSON.stringify({
orderId,
customerId: event.customerId,
items: event.items,
totalAmount: event.totalAmount,
timestamp: Date.now()
})
}]
})
);
return { orderId, status: 'PENDING' };
};
// Inventory Service: OrderCreated'i dinler, stok rezerve eder
export const reserveInventoryOnOrderHandler = async (event: any) => {
const { orderId, items } = event.detail;
try {
const reservationResult = await reserveInventoryForOrder(orderId, items);
const eventBridge = new EventBridgeClient({});
await eventBridge.send(
new PutEventsCommand({
Entries: [{
Source: 'inventory.service',
DetailType: 'InventoryReserved',
Detail: JSON.stringify({
orderId,
reservationId: reservationResult.reservationId,
items,
timestamp: Date.now()
})
}]
})
);
} catch (error: any) {
// Failure event'i publish et
const eventBridge = new EventBridgeClient({});
await eventBridge.send(
new PutEventsCommand({
Entries: [{
Source: 'inventory.service',
DetailType: 'InventoryReservationFailed',
Detail: JSON.stringify({
orderId,
reason: error.message,
timestamp: Date.now()
})
}]
})
);
}
};
// Payment Service: InventoryReserved'i dinler, ödemeyi işler
export const processPaymentOnInventoryReservedHandler = async (event: any) => {
const { orderId } = event.detail;
try {
const order = await getOrder(orderId);
const paymentResult = await processPayment({
orderId,
customerId: order.customerId,
amount: order.totalAmount
});
const eventBridge = new EventBridgeClient({});
await eventBridge.send(
new PutEventsCommand({
Entries: [{
Source: 'payment.service',
DetailType: 'PaymentCompleted',
Detail: JSON.stringify({
orderId,
paymentId: paymentResult.paymentId,
transactionId: paymentResult.transactionId,
timestamp: Date.now()
})
}]
})
);
} catch (error: any) {
// Failure event - compensation'ı trigger eder
const eventBridge = new EventBridgeClient({});
await eventBridge.send(
new PutEventsCommand({
Entries: [{
Source: 'payment.service',
DetailType: 'PaymentFailed',
Detail: JSON.stringify({
orderId,
reason: error.message,
timestamp: Date.now()
})
}]
})
);
}
};
// Inventory Service: PaymentFailed'i dinler, stoğu serbest bırakır (compensation)
export const releaseInventoryOnPaymentFailedHandler = async (event: any) => {
const { orderId } = event.detail;
console.log('Compensating: Ödeme başarısız, stok serbest bırakılıyor', { orderId });
await releaseInventory(orderId);
const eventBridge = new EventBridgeClient({});
await eventBridge.send(
new PutEventsCommand({
Entries: [{
Source: 'inventory.service',
DetailType: 'InventoryReleased',
Detail: JSON.stringify({
orderId,
reason: 'PaymentFailed',
timestamp: Date.now()
})
}]
})
);
};
Choreography’de her servis ilgili event’leri dinlemek ve yeni event’ler publish etmekten sorumlu. Compensation da aynı event mekanizması üzerinden gerçekleşiyor; bir servis failure event’i publish ettiğinde, diğer servisler compensating transaction’larını execute ederek react ediyor.
Semantic Locking ile Isolation#
Saga’lar geleneksel transaction isolation’dan yoksun, bu da concurrent saga conflict’lerine yol açabilir. Semantic locking application-level isolation sağlıyor:
interface SagaLock {
sagaId: string;
lockedAt: number;
expiresAt: number;
}
export const acquireSagaLock = async (orderId: string, sagaId: string) => {
const lockDuration = 5 * 60 * 1000; // 5 dakika
const now = Date.now();
try {
await docClient.send(
new UpdateCommand({
TableName: process.env.ORDERS_TABLE!,
Key: { orderId },
UpdateExpression:
'SET sagaLock = :lock, #status = :processing',
ConditionExpression:
'attribute_not_exists(sagaLock) OR sagaLock.expiresAt < :now',
ExpressionAttributeNames: {
'#status': 'status'
},
ExpressionAttributeValues: {
':lock': {
sagaId,
lockedAt: now,
expiresAt: now + lockDuration
},
':processing': 'PROCESSING',
':now': now
}
})
);
console.log('Saga lock acquired', { orderId, sagaId });
return true;
} catch (error: any) {
if (error.name === 'ConditionalCheckFailedException') {
console.warn('Saga lock already held by another saga', { orderId });
throw new Error('SagaLockConflict');
}
throw error;
}
};
export const releaseSagaLock = async (orderId: string, sagaId: string) => {
await docClient.send(
new UpdateCommand({
TableName: process.env.ORDERS_TABLE!,
Key: { orderId },
UpdateExpression: 'REMOVE sagaLock',
ConditionExpression: 'sagaLock.sagaId = :sagaId',
ExpressionAttributeValues: {
':sagaId': sagaId
}
})
);
console.log('Saga lock released', { orderId, sagaId });
};
Bu semantic lock, iki saga’nın aynı siparişi concurrent olarak modify etmesini önlüyor. Lock, saga crash olduğunda lock’u release edememe durumlarını handle etmek için expiration time içeriyor.
Maliyet Analizi ve Trade-off’lar#
Maliyet etkilerini anlamak doğru yaklaşımı seçmene yardımcı oluyor.
Step Functions Orchestration Maliyetleri#
Ayda 100,000 sipariş için yukarıdaki dört adımlı saga’nın başarılı bir çalışması 5 state’e giriyor: dört task ve terminal Succeed. Catch dalları ancak bir hata oradan geçtiğinde faturalanıyor (us-east-1 liste fiyatları):
- Total state transition: 500,000
- Maliyet: (500,000 / 1,000) × $0.025 = $12.50/ay
- Başarısız saga’lar (%5, her biri 4 compensation transition): 20,000 transition = +$0.50/ay
- Toplam: ~$13/ay
EventBridge Choreography Maliyetleri#
Aynı hacim, sipariş başına 4 event:
- Total event: 400,000
- Maliyet: (400,000 / 1,000,000) × $1.00 = $0.40/ay
- Başarısız siparişler (%5, her biri 2 compensation event): 10,000 event = +$0.01/ay
- Toplam: ~$0.41/ay
Maliyete Karşı Görünürlük#
Bu hacimde choreography yaklaşık otuz kat ucuza geliyor ve aradaki farkı development ile debugging süresinden geri alıyor. Orchestration daha pahalı; karşılığında execution history, hazır retry mekanizması ve “bu sipariş nerede durdu?” sorusunun tek bir cevap adresi geliyor. Çoğu production sisteminde bu görünürlük aylık on iki dolar farka değer; karşılaştırman gereken kalem faturadaki satır değil, bir mühendisin yarım günü.
Yaygın Hatalar ve Çözümler#
Hata 1: Non-Idempotent Operation’lar#
Problem: Retry’da ödeme birden fazla kez alınıyor.
Çözüm: Her zaman idempotency check’leri implement et ve provider idempotency key’leri kullan.
Hata 2: Eksik Compensation Chain’leri#
Problem: Sadece son step compensate ediliyor, daha önceki step’ler inconsistent state’te kalıyor.
Çözüm: Tüm compensation’ları reverse order’da chain’le. Her catch block önceki tüm step’leri compensate etmeli.
Hata 3: Compensation Failure’larını Ignore Etmek#
Problem: Compensation başarısız oluyor, saga askıda kalıyor.
Çözüm: Compensation’lar için aggressive retry (10+ attempt) ve manual intervention için dead-letter queue.
Hata 4: Timeout Çok Kısa#
Problem: Timeout compensation’ı tetikliyor ama operation aslında başarılı olmuştu.
Çözüm: Buffer ile gerçekçi timeout’lar belirle. Compensate etmeden önce gerçek state’i doğrula.
Hata 5: Choreography’de Saga State Takibi Yok#
Problem: Hangi siparişlerin compensation’da olduğu belirlenemiyor.
Çözüm: Observability için choreography’de bile saga state’ini persist et.
Yardımcı Fonksiyonlar#
// Örneklerde referans verilen basitleştirilmiş helper fonksiyonları
import { PutCommand } from '@aws-sdk/lib-dynamodb';
async function saveOrder(order: any) {
await docClient.send(
new PutCommand({
TableName: process.env.ORDERS_TABLE!,
Item: order
})
);
}
async function getOrder(orderId: string) {
const result = await docClient.send(
new GetCommand({
TableName: process.env.ORDERS_TABLE!,
Key: { orderId }
})
);
return result.Item as any;
}
async function reserveInventoryForOrder(orderId: string, items: any[]) {
return { reservationId: `res_${Date.now()}` };
}
async function releaseInventory(orderId: string) {
// Release logic
}
async function processPayment(params: any) {
return {
paymentId: `pay_${Date.now()}`,
transactionId: `txn_${Date.now()}`
};
}
async function updateOrderStatus(orderId: string, status: string) {
await docClient.send(
new UpdateCommand({
TableName: process.env.ORDERS_TABLE!,
Key: { orderId },
UpdateExpression: 'SET #status = :status',
ExpressionAttributeNames: { '#status': 'status' },
ExpressionAttributeValues: { ':status': status }
})
);
}
Hangi Yaklaşım Nerede Geçerli#
Akış üç-dört servisten fazlasına yayılıyorsa, human approval adımı içeriyorsa veya bir incident sırasında birine anlatılması gerekiyorsa varsayılan Step Functions orchestration olsun: execution history hangi step’in başarısız olduğunu ve hangi compensation’ların çalıştığını zaten söylüyor. Akış linear’sa, servisler zaten domain event publish ediyorsa ve merkezi koordinatörün maliyeti (bir deployment ve bir sahip daha) sağladığı görünürlükten ağır basıyorsa EventBridge choreography’ye geç.
İki kısıt her iki seçimde de değişmiyor. Her step’in bir idempotency key’i olmalı, her compensation iki kez çalışmaya dayanmalı; bunlardan birini atlayan saga yalnızca retry sırasında bozulur ve bu, sorunu fark etmek için en kötü andır. Happy path çalıştıktan sonra sıradaki iş, başarısız compensation’lar için bir dead-letter queue kurmak; o kuyruğu bir insanın okuması gerekiyor.
Kaynaklar#
- Pattern: Saga - microservices.io (yeni sekmede açılır) - Chris Richardson’ın koreografi ve orkestrasyon varyantlarıyla Saga pattern’ini kanonik olarak ele aldığı kaynak
- Saga Design Pattern - Azure Architecture Center (yeni sekmede açılır) - Microsoft’un orkestrasyon, telafi işlemleri ve hata yönetimini kapsayan uygulama rehberi
- Compensating Transaction Pattern - Azure Architecture Center (yeni sekmede açılır) - Saga’ların çalışmasını sağlayan geri alma mekanizmasının ayrıntılı incelemesi
- Sagas - ACM SIGMOD 1987 (yeni sekmede açılır) - Garcia-Molina ve Salem’in uzun süreli işlemler için saga kavramını tanıttığı özgün makale
- Saga Pattern - AWS Prescriptive Guidance (yeni sekmede açılır) - Saga pattern’ini mikro hizmet veri kalıcılığına uygulamak için AWS rehberliği
- AWS Step Functions Fiyatlandırması (yeni sekmede açılır) - Standard Workflow’larda state transition başına fiyatlandırma; yukarıdaki orchestration maliyet modelinin dayanağı
- Amazon EventBridge Fiyatlandırması (yeni sekmede açılır) - Event bus’a publish edilen custom event’ler için milyon başına fiyatlandırma; choreography maliyet modelinde kullanıldı
İlgili yazılar
Step Functions ile production serverless workflow kur: Standard ve Express, Distributed Map, error handling ve CDK örnekleriyle maliyet optimizasyonu.
step-functions · aws-cdk · serverless +4
AppSync subscription'ları yalnızca mutation ile tetiklenir. Downstream BFF olaylarını NONE veri kaynaklı bir mutation'a EventBridge ve CDK ile köprülemeyi inceliyorum.
aws · graphql · serverless +4
Transactional Outbox Pattern'in dağıtık sistemlerdeki dual-write problemini nasıl çözdüğünü, PostgreSQL, DynamoDB ve CDC araçlarıyla pratik implementasyonlarını öğren.
distributed-systems · microservices · event-driven +5
AWS Lambda, API Gateway, DynamoDB ve Step Functions için hızlı geri bildirim ve production güvenilirliği sağlayan kapsamlı bir test stratejisi oluşturmayı öğrenin.
lambda · testing · serverless +8
Builder pattern'in TypeScript tip sistemiyle serverless, veri katmanı ve testte güvenli, keşfedilebilir API'leri nasıl kurduğu, çalışan örneklerle.
typescript · design-patterns · aws-cdk +2