مهندسی تابآوری در سیستمهای تراکنشی توزیعشده: همزمانی و ایدامپوتنسی در داتنت
نویسنده: وحید نصیری
تاریخ: ۱۴۰۵/۰۶/۱۳ ۰۹:۲۵
آدرس: www.dntips.ir
چکیده: در معماری سامانههای مدرن توزیعشده و بهویژه پلتفرمهای پرداخت، اجرای پردازشها تحت عنوان «دقیقاً یکبار» (Exactly-Once Execution) بیشتر یک توهم انتزاعی است تا واقعیتی عملیاتی. وقوع تأخیرهای شبکه (Timeouts)، ارسال مجدد درخواستها توسط کلاینتها (Retries)، جابهجایی بار میان سرورها (Failovers) و بازتحویل پیامها در صفها (Message Redelivery)، مسائلی طبیعی در بستر شبکه هستند. عدم تمهید سازوکارهای مناسب در مواجهه با این رخدادها، منجر به بروز باگهای همزمانی (Concurrency Bugs) نظیر مسابقات بر سر دسترسی به داده (Race Conditions)، بهروزرسانیهای گمشده (Lost Updates) و پردازشهای تکراری با تبعات مالی فاجعهبار میشود. این مقاله به بررسی جامع دو مفهوم بههمپیوسته «کنترل همزمانی» و «ایدامپوتنسی» (همگرایی به نتیجه یکسان یا Idempotency) میپردازد و راهکارهای پیادهسازی سازگار با اکوسیستم Microsoft .NET را با تمرکز بر تعاملات پایگاهداده، مدلسازی داده و الگوهای API ارائه میدهد.
// نمونه کد آسیبپذیر در برابر Race Condition
public void Withdraw(Account account, decimal amount)
{
if (account.Balance >= amount)
{
// در این بازه زمانی، ریسمان دیگری میتواند مقدار موجودی را کاهش دهد
account.Balance -= amount;
}
}account.Balance >= amount را پشت سر میگذارند و مجموعاً ۱۶۰ واحد کسر خواهد شد؛ وضعیتی که به Race Condition مشهور است.UPDATE Accounts SET Balance = Balance - @Amount WHERE Id = @Id AND Balance >= @Amount;
public class Account
{
public Guid Id { get; set; }
public decimal Balance { get; set; }
[Timestamp]
public byte[] RowVersion { get; set; } = default!;
}DbUpdateConcurrencyException را صادر میکند و اپلیکیشن میتواند با بازخوانی وضعیت جدید، سیاست تلاش مجدد یا اعلام خطا را در پیش بگیرد.Interlocked.AddSemaphoreSlimConcurrentDictionary و Channel به جای کالکشنهای پایه.// رویکرد شکننده در مواجهه با دو درخواست همزمان var existing = await _store.GetResponseAsync(idempotencyKey); if (existing != null) return existing; var result = await ProcessChargeAsync(request); await _store.SaveResponseAsync(idempotencyKey, result); return result;
ProcessChargeAsync دو بار اجرا میشود.[درخواست جدید]
│
▼
درج کلید با وضعیت Pending (با Unique Index)
│
├─► درج موفق بود ──► انجام عملیات مالی ──► ثبت وضعیت Completed و ذخیره نتیجه
│
└─► خطای یکتایی (کلید وجود دارد)
│
├─► وضعیت Pending است ──► انتظار کوتاه یا بازگرداندن 409 Conflict
└─► وضعیت Completed است ──► بازگرداندن نتیجه ذخیرهشده اولیهpublic class IdempotencyRecord
{
public string Key { get; set; } = string.Empty;
public string UserId { get; set; } = string.Empty;
public string Path { get; set; } = string.Empty;
public string Status { get; set; } = "Pending"; // Pending, Completed
public string? ResponseBody { get; set; }
public int? StatusCode { get; set; }
public DateTime CreatedAtUtc { get; set; } = DateTime.UtcNow;
}
public async Task<IActionResult> ProcessPaymentAsync(
[FromHeader(Name = "X-Idempotency-Key")] string idempotencyKey,
[FromBody] PaymentRequest request)
{
var userId = GetCurrentUserId();
var record = new IdempotencyRecord
{
Key = idempotencyKey,
UserId = userId,
Path = HttpContext.Request.Path,
Status = "Pending"
};
try
{
// ردیف ابتدایی با کلید یکتا درج میشود
await _dbContext.IdempotencyRecords.AddAsync(record);
await _dbContext.SaveChangesAsync();
}
catch (DbUpdateException) // بروز خطا در شاخص Unique Index
{
var existing = await _dbContext.IdempotencyRecords
.AsNoTracking()
.FirstOrDefaultAsync(r => r.Key == idempotencyKey && r.UserId == userId);
if (existing?.Status == "Completed")
{
return StatusCode(existing.StatusCode!.Value, existing.ResponseBody);
}
// اگر عملیات قبلی هنوز در حال اجرا باشد
return Conflict("عملیات درخواستی در حال حاضر در حال پردازش است.");
}
// فراخوانی سرویس پرداخت با در نظر گرفتن خطا
var result = await _paymentGateway.ChargeAsync(request);
record.Status = "Completed";
record.StatusCode = StatusCodes.Status200OK;
record.ResponseBody = JsonSerializer.Serialize(result);
await _dbContext.SaveChangesAsync();
return Ok(result);
}| لایه سیستم | راهکار پیادهسازی ایدامپوتنسی |
| لایه وب (API Layer) | استفاده از سربرگ کلید ایدامپوتنسی (X-Idempotency-Key) و بازپخش کش پاسخها. |
| پایگاهداده (Database) | استفاده از قیود یکتایی (UNIQUE Constraints)، دستورهای UPSERT و اجرای متدهای شرطی. |
| صف پیام (Message Queue) | الگوی تحویل حداقل یکبار (At-Least-Once Delivery) در Kafka یا RabbitMQ نیازمند ثبت شناسه پیامهای پردازششده در یک جدول میانی همراه با تراکنش محلی (Outbox/Inbox Pattern) است. |
| یکپارچهسازی خارجی (Third-Party) | ارسال کلیدهای ایدامپوتنسی بومی به درگاههای پرداخت بالادستی جهت ممانعت از برداشتهای تکراری در هسته بانکی. |
[Fact]
public async Task ConcurrentWithdrawals_ShouldNotViolateBalanceConstraint()
{
// چیدمان محیط تست
var accountId = Guid.NewGuid();
await InitializeAccountWithBalance(accountId, initialBalance: 100m);
int concurrencyLevel = 10;
decimal withdrawAmount = 20m;
// اجرای همزمان 10 درخواست برداشت روی حسابی با 100 واحد موجودی
var tasks = Enumerable.Range(0, concurrencyLevel)
.Select(_ => _accountService.WithdrawAsync(accountId, withdrawAmount));
await Task.WhenAll(tasks);
// بررسی انطباق موجودی
var finalBalance = await GetAccountBalance(accountId);
// موجودی نباید منفی شود؛ حداکثر 5 تراکنش مجاز به تکمیل بودهاند
Assert.True(finalBalance >= 0);
Assert.Equal(0m, finalBalance);
}در سامانههای تراکنشی توزیعشده، مفهوم «دقیقاً یکبار» (Exactly-Once) وجود خارجی ندارد. تأخیر شبکه (Network Timeout)، قطع ناگهانی سوکت یا اجرای مجدد خودکار توسط کلاینتها (Client Retries)، همگی سناریوهای معمول در محیط عملیاتی هستند. همانطور که در معماری Stripe و الگوی مرجع Rocket Rides (ارائهشده توسط Brandur Leach) تبیین شده، چالش اصلی زمانی رخ میدهد که تراکنش در میانه راه قطع شود؛ یعنی سرویس بیرونی (مانند درگاه پرداخت) کارت را شارژ کرده، اما پاسخ به سرور نرسیده یا قبل از ذخیره در دیتابیس داخلی، کانتینر یا پردازش سرور سقوط (Crash) کرده است. در این مقاله، به پیادهسازی گامبهگام و علمی این الگو در اکوسیستم .NET 9/10 با اتکا به Entity Framework Core، پایگاهداده PostgreSQL و متدهای مدیریت تراکنش، همزمانی و قفلگذاری سطرها میپردازیم.
Ride یا Order).public enum RecoveryPoint
{
Started = 0,
RideCreated = 1,
ChargeCompleted = 2,
Finished = 3
}
public class IdempotencyKeyRecord
{
public long Id { get; set; }
// کلیدها همیشه باید مقید به شناسه کاربر/سازمان باشند
public string UserId { get; set; } = string.Empty;
public string Key { get; set; } = string.Empty;
// هش بدنه و پارامترهای درخواست برای رد کردن درخواستهای متناقض
public string RequestHash { get; set; } = string.Empty;
// زمان قفل برای مدیریت Lease
public DateTimeOffset? LockedAt { get; set; }
// مرحله بازیابی
public RecoveryPoint RecoveryPoint { get; set; } = RecoveryPoint.Started;
public int? ResponseStatusCode { get; set; }
public string? ResponseBody { get; set; }
public DateTimeOffset CreatedAt { get; set; } = DateTimeOffset.UtcNow;
}
public class Ride
{
public long Id { get; set; }
public string UserId { get; set; } = string.Empty;
public long IdempotencyKeyId { get; set; }
public int AmountCents { get; set; }
public string? ChargeId { get; set; }
public DateTimeOffset CreatedAt { get; set; } = DateTimeOffset.UtcNow;
public IdempotencyKeyRecord IdempotencyKey { get; set; } = null!;
}public class AppDbContext : DbContext
{
public DbSet<IdempotencyKeyRecord> IdempotencyKeys => Set<IdempotencyKeyRecord>();
public DbSet<Ride> Rides => Set<Ride>();
public AppDbContext(DbContextOptions<AppDbContext> options) : base(options) { }
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
base.OnModelCreating(modelBuilder);
modelBuilder.Entity<IdempotencyKeyRecord>(b =>
{
b.ToTable("idempotency_keys");
b.HasKey(x => x.Id);
b.Property(x => x.Key).HasMaxLength(255).IsRequired();
b.Property(x => x.UserId).HasMaxLength(128).IsRequired();
b.Property(x => x.RequestHash).HasMaxLength(64).IsRequired();
// قید یکتایی مرکب: کلید به ازای هر کاربر یکتاست
b.HasIndex(x => new { x.UserId, x.Key }).IsUnique();
});
modelBuilder.Entity<Ride>(b =>
{
b.ToTable("rides");
b.HasKey(x => x.Id);
// قید حیاتی ایناریانت دیتابیس: هر کلید حداکثر یک رکورد سفر میسازد
b.HasIndex(x => x.IdempotencyKeyId).IsUnique();
b.HasOne(x => x.IdempotencyKey)
.WithMany()
.HasForeignKey(x => x.IdempotencyKeyId)
.OnDelete(DeleteBehavior.Restrict);
});
}
}Replayed: true).409 Conflict بازگردانده میشود (مگر اینکه زمان اجاره قفل یا Lease منقضی شده باشد که در آن صورت درخواست جدید کار را به دست میگیرد).public static class RequestHasher
{
public static string ComputeSha256(string method, string path, object body)
{
// نرمالسازی بدنه به فرم استاندارد
var json = JsonSerializer.Serialize(body, new JsonSerializerOptions
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
WriteIndented = false
});
var raw = $"{method.ToUpperInvariant()}:{path.ToLowerInvariant()}:{json}";
var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(raw));
return Convert.ToHexString(bytes);
}
}FOR UPDATE و clock_timestamp()، ترکیب EF Core با دستورات خام SQL کارآمدترین شیوه است:public record ClaimResult(
bool Success,
int StatusCode,
string? ResponseBody = null,
IdempotencyKeyRecord? Record = null);
public async Task<ClaimResult> ClaimKeyAsync(
string userId,
string key,
string requestHash,
TimeSpan leaseTimeout,
CancellationToken ct)
{
// ۱. درج کلید در صورت عدم وجود (Insert ... ON CONFLICT DO NOTHING)
await _dbContext.Database.ExecuteSqlInterpolatedAsync($@"
INSERT INTO idempotency_keys (user_id, key, request_hash, recovery_point, created_at)
VALUES ({userId}, {key}, {requestHash}, {RecoveryPoint.Started}, clock_timestamp())
ON CONFLICT (user_id, key) DO NOTHING;", ct);
// ۲. خواندن و قفل سطر با استفاده از SELECT FOR UPDATE
var record = await _dbContext.IdempotencyKeys
.FromSqlInterpolated($@"
SELECT * FROM idempotency_keys
WHERE user_id = {userId} AND key = {key}
FOR UPDATE")
.AsTracking()
.SingleOrDefaultAsync(ct);
if (record == null)
{
return new ClaimResult(false, StatusCodes.Status500InternalServerError, "خطا در بازیابی رکورد کلید.");
}
// ۳. بررسی تغییر پارامترها با همان کلید قبلی
if (!string.Equals(record.RequestHash, requestHash, StringComparison.OrdinalIgnoreCase))
{
return new ClaimResult(false, StatusCodes.Status409Conflict,
"از این Idempotency-Key قبلاً با پارامترهای متفاوتی استفاده شده است.");
}
// ۴. اگر فرایند قبلاً به پایان رسیده، بازپخش پاسخ ذخیرهشده
if (record.ResponseStatusCode.HasValue)
{
return new ClaimResult(true, record.ResponseStatusCode.Value, record.ResponseBody, record);
}
// ۵. بررسی Lease: اگر پردازش در جریان است و منقضی نشده، کلاینت باید بعداً Retry کند
var now = DateTimeOffset.UtcNow;
if (record.LockedAt.HasValue && (now - record.LockedAt.Value) < leaseTimeout)
{
return new ClaimResult(false, StatusCodes.Status409Conflict,
"درخواست دیگری با این کلید هماکنون در حال اجراست.");
}
// ۶. تصاحب قفل (Acquire Lease)
record.LockedAt = now;
await _dbContext.SaveChangesAsync(ct);
return new ClaimResult(true, StatusCodes.Status200OK, null, record);
}[درخواست کلاینت]
│
▼
[فاز ۱: تصاحب کلید و ارزیابی هش]
│
▼
[فاز ۲: ثبت موجودیت محلی (Ride) درون تراکنش EF Core] ──► ذخیره RecoveryPoint = RideCreated
│
▼
[فراخوانی درگاه پرداخت (Stripe) با کلید مشتقشده] ◄── مرز ایزولاسیون بیرونی
│
▼
[فاز ۳: ثبت شناسه پرداخت و تکمیل] ──► ذخیره RecoveryPoint = Finished و آزادسازی LockedAtpublic async Task<IResult> ProcessRidePaymentAsync(
string userId,
string idempotencyKey,
RideRequest request,
CancellationToken ct)
{
var hash = RequestHasher.ComputeSha256("POST", "/api/rides", request);
var leaseTtl = TimeSpan.FromSeconds(15);
// ۱. تصاحب قفل
var claim = await ClaimKeyAsync(userId, idempotencyKey, hash, leaseTtl, ct);
if (!claim.Success)
{
return Results.Json(new { error = claim.ResponseBody }, statusCode: claim.StatusCode);
}
// اگر نتیجه قبلاً ذخیره شده، بازپخش نتیجه همراه با Header اختصاصی
if (claim.ResponseBody != null)
{
return Results.Content(
claim.ResponseBody,
contentType: "application/json",
statusCode: claim.StatusCode);
}
var record = claim.Record!;
// ۲. فاز دو: ثبت موجودیت در دیتابیس داخلی
if (record.RecoveryPoint == RecoveryPoint.Started)
{
await using var tx = await _dbContext.Database.BeginTransactionAsync(ct);
try
{
var ride = new Ride
{
UserId = userId,
IdempotencyKeyId = record.Id,
AmountCents = request.AmountCents
};
_dbContext.Rides.Add(ride);
record.RecoveryPoint = RecoveryPoint.RideCreated;
await _dbContext.SaveChangesAsync(ct);
await tx.CommitAsync(ct);
}
catch (DbUpdateException)
{
await tx.RollbackAsync(ct);
// در صورتی که قفل در اثر تأخیر منقضی شده و درخواست دیگری زودتر ثبت کرده باشد،
// قید Unique Index فعال شده و از ثبت داده اشتباه جلوگیری میکند.
throw;
}
}
// ۳. فراخوانی سامانه خارجی (Third-party) با استفاده از Downstream Idempotency Key
// کلید ارسالی به درگاه باید مشتقشده از کلید اصلی باشد
string downstreamKey = $"{userId}:{idempotencyKey}:charge";
string? chargeId = null;
if (record.RecoveryPoint == RecoveryPoint.RideCreated)
{
// فراخوانی درگاه بانکی
chargeId = await _paymentGateway.CreateChargeAsync(
request.AmountCents,
downstreamKey,
ct);
// شبیهسازی سقوط سرور: اگر در این نقطه سرور کرش کند،
// در تلاش مجدد، برنامه از همین فاز کار را ادامه میدهد
// و همان chargeId را از درگاه مجدداً دریافت میکند.
}
// ۴. فاز سه: تکمیل و ذخیره پاسخ نهایی
await using (var tx = await _dbContext.Database.BeginTransactionAsync(ct))
{
var ride = await _dbContext.Rides
.SingleAsync(r => r.IdempotencyKeyId == record.Id, ct);
ride.ChargeId = chargeId;
var responsePayload = JsonSerializer.Serialize(new
{
rideId = ride.Id,
amount = ride.AmountCents,
chargeId = ride.ChargeId
});
record.RecoveryPoint = RecoveryPoint.Finished;
record.ResponseStatusCode = StatusCodes.Status201Created;
record.ResponseBody = responsePayload;
record.LockedAt = null; // آزادسازی قفل
await _dbContext.SaveChangesAsync(ct);
await tx.CommitAsync(ct);
return Results.Created($"/api/rides/{ride.Id}", JsonSerializer.Deserialize<object>(responsePayload));
}
}CREATE UNIQUE INDEX rides_one_per_key ON rides (idempotency_key_id) که در OnModelCreating نوشتیم ضامن نهایی است. حتی اگر Lease اشتباه کار کند، پایگاهداده جلوی ردیف دوم را با DbUpdateException میگیرد؛ یعنی شکست آشکار (500) رخ میدهد اما پول یا خدمات چندباره داده نمیشود.xmin در Postgres یا فیلد ورژن در SQL Server) باعث میشود رکوردی که Lease آن منقضی شده، در زمان Commit متوجه سلب مالکیت خود بشود.idempotency_keys با گذر زمان حجیم میشود. نگهداری کلیدها بین ۲۴ تا ۷۲ ساعت معمول است. در داتنت میتوان از IHostedService یا BackgroundService برای پاکسازی دورهای استفاده کرد:public class IdempotencyReaperWorker : BackgroundService
{
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<IdempotencyReaperWorker> _logger;
public IdempotencyReaperWorker(
IServiceProvider serviceProvider,
ILogger<IdempotencyReaperWorker> logger)
{
_serviceProvider = serviceProvider;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
try
{
using var scope = _serviceProvider.CreateScope();
var dbContext = scope.ServiceProvider.GetRequiredService<AppDbContext>();
var retentionThreshold = DateTimeOffset.UtcNow.AddHours(-48);
// استفاده از ExecuteDeleteAsync در EF Core 7/8/9 بدون بارگذاری در حافظه
int deletedCount = await dbContext.IdempotencyKeys
.Where(k => k.CreatedAt < retentionThreshold && k.ResponseStatusCode != null)
.ExecuteDeleteAsync(stoppingToken);
if (deletedCount > 0)
{
_logger.LogInformation("تعداد {Count} کلید ایدامپوتنسی قدیمی پاکسازی شدند.", deletedCount);
}
}
catch (Exception ex)
{
_logger.LogError(ex, "خطا حین اجرای پاکسازی دورهای کلیدهای ایدامپوتنسی.");
}
// اجرای هر ۴ ساعت یکبار
await Task.Delay(TimeSpan.FromHours(4), stoppingToken);
}
}
}(UserId, Path, Key) یکتا کنید.key:charge) بفرستید.Unique Constraint).BeginTransactionAsync تفکیکشده به ازای هر گام استفاده کنید و وضعیت RecoveryPoint را پیش از خروج از تراکنش ثبت کنید.Idempotent-Replayed: true را در خط لوله Middleware برگردانید تا تیم فرانتاند یا سیستمهای مصرفکننده متوجه کشف مجدد پاسخ شوند.