如何使用C# MongoDB强类型驱动实现两阶段提交?
Alright, let's tackle this two-phase commit implementation for your MongoDB + C# driver 2.5.0 scenario. Since you need all-or-nothing updates across multiple collections, following MongoDB's two-phase commit pattern is the right call—here's a step-by-step guide tailored to your setup.
Before diving into code, let's make sure we're on the same page:
- Prepare Phase: Verify all intended updates are possible, mark documents as part of a pending transaction, and create a transaction record to track progress.
- Commit Phase: Apply all actual updates, then mark the transaction as completed.
- Rollback Phase: If anything fails during prepare or commit, undo all partial changes and mark the transaction as rolled back.
- You'll need a dedicated
transactionscollection to track every transaction's state—this is critical for recovery if something goes wrong mid-process.
2.1 Define Models for Transaction Tracking
First, create strongly-typed classes to represent your transaction records and the operations they contain. These will map directly to documents in the transactions collection:
using MongoDB.Bson; using MongoDB.Bson.Serialization.Attributes; public enum TransactionState { Prepared, Committed, RolledBack, Completed } public class Transaction { [BsonId] public ObjectId Id { get; set; } [BsonElement("state")] public TransactionState State { get; set; } [BsonElement("operations")] public List<TransactionOperation> Operations { get; set; } [BsonElement("timestamp")] public DateTime Timestamp { get; set; } [BsonElement("retryCount")] public int RetryCount { get; set; } } public class TransactionOperation { [BsonElement("collection")] public string CollectionName { get; set; } [BsonElement("filter")] public BsonDocument Filter { get; set; } [BsonElement("update")] public BsonDocument Update { get; set; } [BsonElement("isApplied")] public bool IsApplied { get; set; } }
2.2 Implement the Prepare Phase
This phase ensures all updates can be executed before we commit anything. We'll add a temporary marker to target documents to track they're part of a pending transaction:
public async Task<ObjectId> PrepareTransaction(IMongoDatabase db, List<TransactionOperation> operations) { var transactionCollection = db.GetCollection<Transaction>("transactions"); var transactionId = ObjectId.GenerateNewId(); var transaction = new Transaction { Id = transactionId, State = TransactionState.Prepared, Operations = operations.Select(op => new TransactionOperation { CollectionName = op.CollectionName, Filter = op.Filter, Update = op.Update, IsApplied = false }).ToList(), Timestamp = DateTime.UtcNow, RetryCount = 0 }; try { // Verify each document exists and mark it as part of this transaction foreach (var op in transaction.Operations) { var targetCollection = db.GetCollection<BsonDocument>(op.CollectionName); var result = await targetCollection.FindOneAndUpdateAsync( op.Filter, Builders<BsonDocument>.Update.Set("__pendingTransactionId", transactionId), new FindOneAndUpdateOptions<BsonDocument> { ReturnDocument = ReturnDocument.After }); if (result == null) { // Document not found—rollback any already marked documents await RollbackPreparedOperations(db, transaction); throw new Exception($"Document matching filter {op.Filter} not found in {op.CollectionName}"); } op.IsApplied = true; } // Save the prepared transaction record await transactionCollection.InsertOneAsync(transaction); return transactionId; } catch (Exception ex) { await RollbackPreparedOperations(db, transaction); throw new Exception("Prepare phase failed", ex); } } private async Task RollbackPreparedOperations(IMongoDatabase db, Transaction transaction) { foreach (var op in transaction.Operations.Where(o => o.IsApplied)) { var targetCollection = db.GetCollection<BsonDocument>(op.CollectionName); // Remove the temporary transaction marker await targetCollection.UpdateManyAsync( op.Filter, Builders<BsonDocument>.Update.Unset("__pendingTransactionId")); } }
2.3 Implement the Commit Phase
Once the prepare phase succeeds, we can apply the actual updates and clean up the temporary markers:
public async Task CommitTransaction(IMongoDatabase db, ObjectId transactionId) { var transactionCollection = db.GetCollection<Transaction>("transactions"); var transaction = await transactionCollection .Find(t => t.Id == transactionId && t.State == TransactionState.Prepared) .FirstOrDefaultAsync(); if (transaction == null) { throw new Exception("Transaction not found or already processed"); } try { foreach (var op in transaction.Operations) { var targetCollection = db.GetCollection<BsonDocument>(op.CollectionName); // Merge the actual update with a cleanup of the temporary marker var combinedUpdate = op.Update.Merge( Builders<BsonDocument>.Update.Unset("__pendingTransactionId")); var result = await targetCollection.UpdateManyAsync(op.Filter, combinedUpdate); if (result.ModifiedCount == 0) { throw new Exception($"Failed to apply update to {op.CollectionName} (no documents modified)"); } } // Mark the transaction as successfully committed await transactionCollection.UpdateOneAsync( t => t.Id == transactionId, Builders<Transaction>.Update.Set(t => t.State, TransactionState.Committed)); } catch (Exception ex) { // Increment retry count for later recovery attempts transaction.RetryCount++; await transactionCollection.ReplaceOneAsync(t => t.Id == transactionId, transaction); throw new Exception("Commit phase failed", ex); } }
2.4 Implement the Rollback Phase
If the commit phase fails (or a transaction times out), we need to undo all partial changes:
public async Task RollbackTransaction(IMongoDatabase db, ObjectId transactionId) { var transactionCollection = db.GetCollection<Transaction>("transactions"); var transaction = await transactionCollection .Find(t => t.Id == transactionId && t.State == TransactionState.Prepared) .FirstOrDefaultAsync(); if (transaction == null) return; // Transaction already processed try { foreach (var op in transaction.Operations.Where(o => o.IsApplied)) { var targetCollection = db.GetCollection<BsonDocument>(op.CollectionName); await targetCollection.UpdateManyAsync( op.Filter, Builders<BsonDocument>.Update.Unset("__pendingTransactionId")); } // Mark the transaction as rolled back await transactionCollection.UpdateOneAsync( t => t.Id == transactionId, Builders<Transaction>.Update.Set(t => t.State, TransactionState.RolledBack)); } catch (Exception ex) { transaction.RetryCount++; await transactionCollection.ReplaceOneAsync(t => t.Id == transactionId, transaction); throw new Exception("Rollback phase failed", ex); } }
2.5 Add Background Recovery Logic
Since MongoDB's two-phase commit isn't fully atomic, you'll need a background task to handle stale or failed transactions:
public async Task ProcessStaleTransactions(IMongoDatabase db, TimeSpan timeout, int maxRetries) { var transactionCollection = db.GetCollection<Transaction>("transactions"); var staleTransactions = await transactionCollection.Find(t => t.State == TransactionState.Prepared && t.Timestamp < DateTime.UtcNow.Subtract(timeout)).ToListAsync(); foreach (var transaction in staleTransactions) { if (transaction.RetryCount < maxRetries) { try { await CommitTransaction(db, transaction.Id); } catch { transaction.RetryCount++; await transactionCollection.ReplaceOneAsync(t => t.Id == transaction.Id, transaction); } } else { await RollbackTransaction(db, transaction.Id); } } // Clean up old completed transactions (keep for 7 days as an example) await transactionCollection.DeleteManyAsync(t => (t.State == TransactionState.Committed || t.State == TransactionState.RolledBack) && t.Timestamp < DateTime.UtcNow.Subtract(TimeSpan.FromDays(7))); }
- Index the Transactions Collection: Add indexes on
StateandTimestampto make the stale transaction scan faster:await transactionCollection.Indexes.CreateOneAsync( Builders<Transaction>.IndexKeys.Ascending(t => t.State).Ascending(t => t.Timestamp)); - Strongly-Typed Collections: If you're using strongly-typed collections instead of
BsonDocument, replace the generic type inGetCollection<T>with your entity class. - Concurrency Safety: Add a
Versionfield to theTransactionclass and use it in update operations to prevent concurrent modification of the same transaction record.
内容的提问来源于stack exchange,提问作者Waxren

