MongoDB C# ChangeStream报错:$Index数组索引越界
MongoDB C# Driver ChangeStream 运行数小时后触发数组越界异常
问题描述
我们使用MongoDB C# Driver的ChangeStream实现实时同步业务,监听多个数据库集合的插入与更新操作,但每次运行数小时后都会触发以下异常:
Index was outside the bounds of the array.
异常堆栈信息
at Microsoft.Extensions.Logging.LogValuesFormatter.GetValue(Object[] values, Int32 index) at Microsoft.Extensions.Logging.FormattedLogValues.GetEnumerator()+MoveNext() at Serilog.Extensions.Logging.SerilogLogger.Log[TState](LogLevel logLevel, EventId eventId, TState state, Exception exception, Func`3 formatter) at Microsoft.Extensions.Logging.Logger`1.Microsoft.Extensions.Logging.ILogger.Log[TState](LogLevel logLevel, EventId eventId, TState state, Exception exception, Func`3 formatter) at Microsoft.Extensions.Logging.LoggerExtensions.Log(ILogger logger, LogLevel logLevel, EventId eventId, Exception exception, String message, Object[] args) at SyncServer.Services.MultiDatabaseChangeStreamManager.<>cDisplayClass9_0.<<StartChangeStreams>b6>d.MoveNext() in /app/jenkins/workspace/flecx-prod-changestream/src/Host/SyncServer/Services/MultiDatabaseChangeStreamManager.cs:line 167 — End of stack trace from previous location — at MongoDB.Driver.IAsyncCursorExtensions.ForEachAsync[TDocument](IAsyncCursor`1 source, Func`3 processor, CancellationToken cancellationToken) at SyncServer.Services.MultiDatabaseChangeStreamManager.StartChangeStreams() in /app/jenkins/workspace/flecx-prod-changestream/src/Host/SyncServer/Services/MultiDatabaseChangeStreamManager.cs:line 109
第109行对应ChangeStream的ForEachAsync执行逻辑,核心代码如下:
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<BsonDocument>>() .Match(change => change.OperationType == ChangeStreamOperationType.Insert || change.OperationType == ChangeStreamOperationType.Update) .Match(y => client.SyncSettings.WatchedDatabases.Contains(y.DatabaseNamespace.DatabaseName)) .Match(x => client.SyncSettings.WatchedCollections.Contains(x.CollectionNamespace.CollectionName)); var resumeToken = new BsonDocument(); var logCollection = database.GetCollection<Logstream>("logstream"); var lastLog = logCollection.AsQueryable().Where(x => x.Status == LogstreamStatus.Completed).OrderByDescending(o => o.CreatedOn).FirstOrDefault(); if (lastLog is not null) { _logger.LogInformation("Last failed log found"); resumeToken = lastLog.ResumeToken; } else { _logger.LogInformation("No last log found, fetching resume token from mongodb directly"); var previousCursor = await database.WatchAsync(); resumeToken = previousCursor.GetResumeToken(); _logger.LogInformation("Resume token fetched"); previousCursor.Dispose(); } var options = new ChangeStreamOptions { ResumeAfter = resumeToken, BatchSize = 5 }; CancellationTokenSource cancelTokenSource = new CancellationTokenSource(); CancellationToken token = cancelTokenSource.Token; try { using (var cursor = await mongoclient.WatchAsync(pipeline, options, cancellationToken: token)) { await cursor.ForEachAsync(async change => { _logger.LogInformation("Inspecting change for | {documentId} | {operationType}", change.DocumentKey["_id"], change.OperationType); var log = Logstream.Create(client.ClientId, change.DocumentKey, change.FullDocument, change.OperationType == ChangeStreamOperationType.Insert ? null : change.UpdateDescription.UpdatedFields, change.ResumeToken, change.OperationType, change.CollectionNamespace.CollectionName); try { switch (change.OperationType) { case ChangeStreamOperationType.Insert: { _logger.LogInformation($"Insert object = {change.FullDocument.ToJson()}"); _logger.LogInformation($"ReplayChangeToTarget start"); await ReplayChangeToTarget(change, client.SyncSettings, change.CollectionNamespace.CollectionName); _logger.LogInformation($"ReplayChangeToTarget end"); _logger.LogInformation($"Inserting logs"); await logCollection.InsertOneAsync(log); _logger.LogInformation($"logs inserted"); break; } case ChangeStreamOperationType.Update: { _logger.LogInformation($"Update object = {change.UpdateDescription.UpdatedFields.ToJson()}"); _logger.LogInformation($"ReplayChangeToTarget start"); await ReplayChangeToTarget(change, client.SyncSettings, change.CollectionNamespace.CollectionName); _logger.LogInformation($"ReplayChangeToTarget end"); _logger.LogInformation($"Inserting logs"); await logCollection.InsertOneAsync(log); _logger.LogInformation($"logs inserted"); break; } } } catch (MongoWriteException ex) { if (ex.WriteError.Code == 11000) { Console.WriteLine($"Exception Dup key: {change.FullDocument.GetValue("_id")}, ignore..."); log.Status = LogstreamStatus.Completed; await logCollection.InsertOneAsync(log); } else { _logger.LogInformation($"MongoException Exception :: {ex.Message}", ex.ToString()); await SendEmailReport(client.ClientName, change.CollectionNamespace.CollectionName, ex.Message); cancelTokenSource.Cancel(); } } catch (Exception ex) { _logger.LogInformation($"ReplayChangeToTarget Exception :: {ex.Message}", ex.ToString()); await SendEmailReport(client.ClientName, change.CollectionNamespace.CollectionName, ex.Message); cancelTokenSource.Cancel(); } }, cancellationToken: token); } } catch (Exception ex) { await SendEmailReport("App", "App", ex.Message); _logger.LogInformation($"Collection.WatchAsync Exception :: ${ex.Message}", ex.ToString()); _logger.LogInformation(ex.StackTrace); } private async Task ReplayChangeToTarget(ChangeStreamDocument<BsonDocument> change, SyncSettings settings, string collectionName) { _logger.LogInformation($"Sync target started for {collectionName}"); IMongoDatabase database; IMongoClient client; if (settings.IsSSLEnable) { var path = Environment.CurrentDirectory + Path.DirectorySeparatorChar + "Certificates" + Path.DirectorySeparatorChar + settings.CertificateFilePath; client = _dbManager.GetClientFromUrl(settings.ConnectionString, path, settings.Password); database = client.GetDatabase(settings.Database); } else { database = _dbManager.GetMongoDbFromUrl(settings.ConnectionString); } var collection = database.GetCollection<BsonDocument>(collectionName); switch (change.OperationType) { case ChangeStreamOperationType.Insert: { _logger.LogInformation($"Inserting to target collection = {collectionName}"); await collection.InsertOneAsync(change.FullDocument); _logger.LogInformation($"Inserted successfully to target collection = {collectionName}"); break; } case ChangeStreamOperationType.Update: { _logger.LogInformation($"Updating to target collection = {collectionName}"); var filter = Builders<BsonDocument>.Filter.Eq("_id", change.DocumentKey["_id"]); var fields = change.UpdateDescription.UpdatedFields; var updateDefination = new List<UpdateDefinition<BsonDocument>>(); foreach (var dataField in fields) { updateDefination.Add(Builders<BsonDocument>.Update.Set(dataField.Name, dataField.Value)); } if (updateDefination.Count > 0) { var combinedUpdate = Builders<BsonDocument>.Update.Combine(updateDefination); await collection.UpdateOneAsync(filter, combinedUpdate); _logger.LogInformation($"Updated successfully to target collection = {collectionName}"); } break; } } _logger.LogInformation($"Sync target ended for {collectionName}"); }
问题分析与解决方案
从堆栈信息看,异常出现在日志格式化环节,并非MongoDB Driver本身问题,而是Microsoft.Extensions.Logging的LogValuesFormatter在处理日志参数时触发数组越界。
触发原因
代码中存在多处日志调用参数不匹配的情况:
_logger.LogInformation($"MongoException Exception :: {ex.Message}", ex.ToString());:格式化字符串仅1个占位符,但传入2个参数,日志组件会尝试访问超出数组长度的索引,运行一段时间后积累的上下文触发异常。_logger.LogInformation($"ReplayChangeToTarget Exception :: {ex.Message}", ex.ToString());:同上,参数数量与占位符不匹配。_logger.LogInformation($"Collection.WatchAsync Exception :: ${ex.Message}", ex.ToString());:格式化字符串写法错误(多余$符号),同时存在参数多传问题。
修复方案
修正所有日志调用的参数匹配问题,优先使用带Exception参数的日志重载记录完整异常:
// 错误写法 // _logger.LogInformation($"MongoException Exception :: {ex.Message}", ex.ToString()); // 正确写法:用LogError重载记录异常详情 _logger.LogError(ex, $"MongoException Exception :: {ex.Message}");
其余类似日志调用需同步修正,确保格式化字符串的占位符数量与传入参数数量一致。
额外优化建议
- 异常自动恢复:当前捕获异常后直接停止监听,建议添加自动重试逻辑,利用
ResumeAfter或StartAfter恢复ChangeStream,避免服务中断。 - 日志级别规范:使用
LogError记录异常,而非LogInformation,便于日志分级排查。 - 异步逻辑校验:
cursor.ForEachAsync内的异步委托需确保无未处理异常,可添加全局异常捕获兜底。
内容的提问来源于stack exchange,提问作者Ghazanfar Khan
相关产品推荐
相关产品推荐

