You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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在处理日志参数时触发数组越界。

触发原因

代码中存在多处日志调用参数不匹配的情况:

  1. _logger.LogInformation($"MongoException Exception :: {ex.Message}", ex.ToString());:格式化字符串仅1个占位符,但传入2个参数,日志组件会尝试访问超出数组长度的索引,运行一段时间后积累的上下文触发异常。
  2. _logger.LogInformation($"ReplayChangeToTarget Exception :: {ex.Message}", ex.ToString());:同上,参数数量与占位符不匹配。
  3. _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}");

其余类似日志调用需同步修正,确保格式化字符串的占位符数量与传入参数数量一致。

额外优化建议

  1. 异常自动恢复:当前捕获异常后直接停止监听,建议添加自动重试逻辑,利用ResumeAfter或StartAfter恢复ChangeStream,避免服务中断。
  2. 日志级别规范:使用LogError记录异常,而非LogInformation,便于日志分级排查。
  3. 异步逻辑校验:cursor.ForEachAsync内的异步委托需确保无未处理异常,可添加全局异常捕获兜底。

内容的提问来源于stack exchange,提问作者Ghazanfar Khan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 20:04:51