using System.Collections.Concurrent;
using System.Diagnostics;
using System.IO.Pipelines;
using System.Security.Cryptography;
using System.Text;
using Dpz.Core.MongodbAccess;
using ICSharpCode.SharpZipLib.Core;
using ICSharpCode.SharpZipLib.Zip;
using Microsoft.Extensions.Logging;
using MongoDB.Bson;
using MongoDB.Driver;

namespace Dpz.Core.Backup;

/// <summary>
/// 数据备份
/// </summary>
public class BackupManager : IDisposable
{
    /// <summary>
    /// 最大失败重试次数
    /// </summary>
    private const int MaxRetries = 15;

    /// <summary>
    /// 重试间隔时间
    /// </summary>
    private readonly TimeSpan _retryDelay = TimeSpan.FromSeconds(5);

    // 检查点数据类
    private record CheckpointData(
        string CollectionName,
        string LastDocumentId,
        string IdType,
        bool IsCompleted
    );

    // 改为实例级别,每次备份都是全新的检查点
    private readonly ConcurrentDictionary<string, CheckpointData> _documentsCheckpoint = new();

    // 活跃备份任务跟踪,防止同一数据库被多个实例同时备份
    private static readonly ConcurrentDictionary<string, string> ActiveBackups = new();

    private readonly string _instanceId = Guid.NewGuid().ToString("N");
    private readonly string _backupKey;
    private readonly IRepositoryBase _repositoryBase;
    private readonly string _backupPath;
    private readonly ILogger<BackupManager> _logger;
    private readonly string _today;
    private readonly string _database;
    private readonly CancellationToken _cancellationToken;

    public BackupManager(
        string? connectionString,
        string backupKey,
        ILoggerFactory loggerFactory,
        CancellationToken cancellationToken = default
    )
    {
        _backupKey = backupKey;
        _repositoryBase = new RepositoryBase(connectionString);
        _logger = loggerFactory.CreateLogger<BackupManager>();
        _database = _repositoryBase.Database.DatabaseNamespace.DatabaseName;
        _cancellationToken = cancellationToken;

        _logger.LogInformation(
            "当前备份数据库:{Database},实例ID:{InstanceId}",
            _database,
            _instanceId
        );

        _today = DateTime.Now.ToString("yyyyMMdd");
        var backupPath = Path.Combine("backup", _database, _today);

        // 检查是否已有实例在备份同一数据库
        var backupTaskKey = $"{_database}_{_today}";
        if (!ActiveBackups.TryAdd(backupTaskKey, _instanceId))
        {
            var existingInstanceId = ActiveBackups.TryGetValue(backupTaskKey, out var existing)
                ? existing
                : "未知";
            throw new InvalidOperationException(
                $"数据库 {_database} 在 {_today} 的备份任务已在进行中,"
                    + $"当前实例ID:{_instanceId},已存在实例ID:{existingInstanceId}"
            );
        }

        try
        {
            _logger.LogInformation("当前备份路径:{BackupPath}", backupPath);

            var directoryInfo = new DirectoryInfo(backupPath);
            if (!directoryInfo.Exists)
            {
                directoryInfo.Create();
                _logger.LogInformation("创建备份目录:{BackupPath}", backupPath);
            }

            _backupPath = directoryInfo.FullName;

            _logger.LogInformation(
                "备份任务已注册:{BackupTaskKey} -> {InstanceId}",
                backupTaskKey,
                _instanceId
            );
        }
        catch
        {
            // 如果构造函数失败,清理已注册的活跃备份记录
            ActiveBackups.TryRemove(backupTaskKey, out _);
            throw;
        }
    }

    /// <summary>
    /// 备份数据库
    /// </summary>
    /// <param name="filter">
    /// 需要过滤的集合
    /// Key :   集合名称   如:Config
    /// Value : 过滤条件   如: { CreateTime: { $gt: ISODate('2024-01-01') } }
    /// </param>
    /// <returns></returns>
    public async Task<BackupResult> BackupAsync(Dictionary<string, string>? filter = null)
    {
        _cancellationToken.ThrowIfCancellationRequested();

        var list = await _repositoryBase.GetAllCollectionAsync();
        _logger.LogInformation("开始备份数据库,获取集合:{@Collections}", list);

        var backupStopwatch = Stopwatch.StartNew();
        await Parallel.ForEachAsync(
            list,
            _cancellationToken,
            async (collectionName, _) =>
            {
                await BackupCollectionAsync(filter, collectionName);
            }
        );
        backupStopwatch.Stop();
        _logger.LogInformation(
            "所有集合备份完成,耗时:{ElapsedMilliseconds} ms",
            backupStopwatch.ElapsedMilliseconds
        );

        var zipPath = Path.Combine(_backupPath, "..", $"{_database}_{_today}.zip");
        _logger.LogInformation("数据获取完毕,准备压缩,压缩文件路径:{ZipPath}", zipPath);

        var password = Sha256Encrypt(_backupKey + _database + _today);

        var compressStopwatch = Stopwatch.StartNew();
        CreateZip(zipPath, password, _backupPath);
        compressStopwatch.Stop();
        _logger.LogInformation(
            "压缩完成,压缩文件路径:{ZipPath},耗时:{ElapsedMilliseconds} ms",
            zipPath,
            compressStopwatch.ElapsedMilliseconds
        );

        // 备份完成,清理活跃备份记录
        CompleteBackup();

        return new BackupResult(
            password,
            zipPath,
            _repositoryBase.Database.DatabaseNamespace.DatabaseName
        );
    }

    private async Task BackupCollectionAsync(
        Dictionary<string, string>? filter,
        string collectionName
    )
    {
        for (var attempt = 0; attempt < MaxRetries; attempt++)
        {
            try
            {
                _cancellationToken.ThrowIfCancellationRequested();

                // 读取检查点
                var checkpoint = LoadCheckpoint(collectionName);
                var lastDocumentId = checkpoint?.LastDocumentId ?? string.Empty;
                var idType = checkpoint?.IdType ?? string.Empty;

                _logger.LogInformation(
                    "开始备份集合 {Collection},检查点信息:LastId={LastId}, IdType={IdType}, IsCompleted={IsCompleted}",
                    collectionName,
                    lastDocumentId,
                    idType,
                    checkpoint?.IsCompleted
                );

                FilterDefinition<BsonDocument>? collectionFilter = null;
                if (
                    filter != null
                    && filter.TryGetValue(collectionName, out var value)
                    && BsonDocument.TryParse(value, out var filterValue)
                )
                {
                    _logger.LogInformation(
                        "backup {Database}, collection {Collection} filter:{@Filter}",
                        _database,
                        collectionName,
                        filterValue
                    );
                    collectionFilter = new BsonDocumentFilterDefinition<BsonDocument>(filterValue);
                }

                // 添加ID过滤条件
                if (!string.IsNullOrEmpty(lastDocumentId))
                {
                    var bsonValue = GetBsonValueFromString(lastDocumentId, idType);
                    _logger.LogDebug(
                        "添加ID过滤条件:LastId={LastId}, IdType={IdType}, BsonValue={BsonValue}",
                        lastDocumentId,
                        idType,
                        bsonValue
                    );

                    var idFilter = Builders<BsonDocument>.Filter.Gt("_id", bsonValue);
                    collectionFilter =
                        collectionFilter == null
                            ? idFilter
                            : Builders<BsonDocument>.Filter.And(collectionFilter, idFilter);
                }

                var backupStopwatch = Stopwatch.StartNew();
                var data = _repositoryBase.SearchForAsync(
                    collectionName,
                    collectionFilter,
                    new FindOptions<BsonDocument>
                    {
                        Sort = Builders<BsonDocument>.Sort.Ascending("_id"),
                    },
                    cancellationToken: _cancellationToken
                );
                var (finalLastId, finalIdType) = await SaveDataFileAsync(
                    data,
                    collectionName,
                    lastDocumentId,
                    idType
                );
                backupStopwatch.Stop();
                _logger.LogInformation(
                    "备份集合:{Collection}完成,耗时:{ElapsedMilliseconds} ms",
                    collectionName,
                    backupStopwatch.ElapsedMilliseconds
                );

                // 更新检查点 - 使用实际处理的最后文档ID
                SaveCheckpoint(new CheckpointData(collectionName, finalLastId, finalIdType, true));

                _logger.LogInformation("集合 {Collection} 备份完成并更新检查点", collectionName);
                return;
            }
            catch (Exception e)
            {
                _logger.LogError(
                    e,
                    "备份集合:{Collection}时出错,尝试次数:{Attempt}/{TotalAttempts}",
                    collectionName,
                    attempt + 1,
                    MaxRetries
                );

                if (attempt == MaxRetries - 1)
                {
                    _logger.LogError(
                        "备份集合 {Collection} 达到最大重试次数,放弃重试",
                        collectionName
                    );
                    return;
                }

                _logger.LogInformation(
                    "等待 {Delay} 秒后重试备份集合 {Collection}",
                    _retryDelay.TotalSeconds,
                    collectionName
                );
                await Task.Delay(_retryDelay, _cancellationToken);
            }
        }
    }

    private async Task<(string FinalLastId, string FinalIdType)> SaveDataFileAsync(
        IAsyncEnumerable<BsonDocument> data,
        string collectionName,
        string lastDocumentId,
        string idType
    )
    {
        _cancellationToken.ThrowIfCancellationRequested();

        var filePath = Path.Combine(_backupPath, collectionName + ".bson");
        var tempFilePath = filePath + ".temp";

        _logger.LogDebug(
            "开始保存集合 {Collection} 的数据,临时文件:{TempFile}",
            collectionName,
            tempFilePath
        );

        var documentCount = 0;
        var currentLastId = lastDocumentId;
        var currentIdType = idType;
        var lastCheckpointTime = DateTime.Now;

        await using (
            var fileStream = new FileStream(
                tempFilePath,
                FileMode.Create,
                FileAccess.Write,
                FileShare.None,
                bufferSize: 1 << 12,
                useAsync: true
            )
        )
        {
            var pipeWriter = PipeWriter.Create(
                fileStream,
                new StreamPipeWriterOptions(leaveOpen: true)
            );

            try
            {
                await foreach (var document in data.WithCancellation(_cancellationToken))
                {
                    var bsonBytes = document.ToBson();
                    var memory = pipeWriter.GetMemory(bsonBytes.Length);
                    bsonBytes.CopyTo(memory);
                    pipeWriter.Advance(bsonBytes.Length);

                    documentCount++;
                    if (document.TryGetValue("_id", out var id))
                    {
                        currentLastId = id?.ToString() ?? "";
                        currentIdType = id?.BsonType.ToString() ?? "";
                    }

                    // 每1000个文档更新一次检查点
                    if (documentCount % 1000 == 0)
                    {
                        var now = DateTime.Now;
                        var elapsed = now - lastCheckpointTime;
                        _logger.LogDebug(
                            "更新检查点:Collection={Collection}, DocumentCount={Count}, LastId={LastId}, IdType={IdType}, 耗时={Elapsed}ms",
                            collectionName,
                            documentCount,
                            currentLastId,
                            currentIdType,
                            elapsed.TotalMilliseconds
                        );

                        SaveCheckpoint(
                            new CheckpointData(collectionName, currentLastId, currentIdType, false)
                        );
                        lastCheckpointTime = now;
                        await Task.Delay(TimeSpan.FromMilliseconds(500), _cancellationToken);
                    }
                }
            }
            finally
            {
                await pipeWriter.FlushAsync(_cancellationToken);
                await pipeWriter.CompleteAsync();
            }
        }

        // 验证数据完整性
        if (documentCount > 0)
        {
            _logger.LogInformation(
                "集合 {Collection} 备份完成,共处理 {Count} 个文档,准备替换文件",
                collectionName,
                documentCount
            );

            // 原子性文件替换 - 直接覆盖,避免删除和移动之间的间隙
            File.Move(tempFilePath, filePath, overwrite: true);
            _logger.LogDebug("将临时文件重命名为正式文件:{File}", filePath);
        }
        else
        {
            _logger.LogWarning("集合 {Collection} 没有数据需要备份,删除临时文件", collectionName);
            // 如果没有数据,删除临时文件
            File.Delete(tempFilePath);
        }
        // 返回最后处理的文档ID信息
        return (currentLastId, currentIdType);
    }

    private CheckpointData? LoadCheckpoint(string collectionName)
    {
        var startTime = DateTime.Now;
        _logger.LogDebug("开始读取检查点:Collection={Collection}", collectionName);

        try
        {
            var checkpointKey = $"{_database}_{collectionName}";
            if (_documentsCheckpoint.TryGetValue(checkpointKey, out var checkpoint))
            {
                var elapsed = DateTime.Now - startTime;
                _logger.LogDebug(
                    "读取到检查点:Collection={Collection}, LastId={LastId}, IdType={IdType}, IsCompleted={IsCompleted}, 耗时={Elapsed}ms",
                    collectionName,
                    checkpoint.LastDocumentId,
                    checkpoint.IdType,
                    checkpoint.IsCompleted,
                    elapsed.TotalMilliseconds
                );
                return checkpoint;
            }
        }
        catch (Exception e)
        {
            _logger.LogError(e, "读取检查点失败:Collection={Collection}", collectionName);
        }

        var totalElapsed = DateTime.Now - startTime;
        _logger.LogDebug(
            "未找到有效的检查点:Collection={Collection}, 耗时={Elapsed}ms",
            collectionName,
            totalElapsed.TotalMilliseconds
        );
        return null;
    }

    private void SaveCheckpoint(CheckpointData checkpoint)
    {
        var startTime = DateTime.Now;
        _logger.LogDebug(
            "开始保存检查点:Collection={Collection}, LastId={LastId}, IdType={IdType}",
            checkpoint.CollectionName,
            checkpoint.LastDocumentId,
            checkpoint.IdType
        );

        try
        {
            // 使用内存操作保存检查点
            var checkpointKey = $"{_database}_{checkpoint.CollectionName}";
            _documentsCheckpoint[checkpointKey] = checkpoint;

            var elapsed = DateTime.Now - startTime;
            _logger.LogDebug(
                "保存检查点完成:Collection={Collection}, LastId={LastId}, IdType={IdType}, IsCompleted={IsCompleted}, 耗时={Elapsed}ms",
                checkpoint.CollectionName,
                checkpoint.LastDocumentId,
                checkpoint.IdType,
                checkpoint.IsCompleted,
                elapsed.TotalMilliseconds
            );
        }
        catch (Exception e)
        {
            var elapsed = DateTime.Now - startTime;
            _logger.LogError(
                e,
                "保存检查点失败:Collection={Collection}, 耗时={Elapsed}ms",
                checkpoint.CollectionName,
                elapsed.TotalMilliseconds
            );
        }
    }

    /// <summary>
    /// 手动完成备份任务,清理活跃备份记录
    /// </summary>
    private void CompleteBackup()
    {
        var backupTaskKey = $"{_database}_{_today}";
        if (ActiveBackups.TryRemove(backupTaskKey, out var removedInstanceId))
        {
            _logger.LogInformation(
                "备份任务已完成,清理记录:{BackupTaskKey} -> {InstanceId}",
                backupTaskKey,
                removedInstanceId
            );
        }
    }

    private bool _disposed;

    public void Dispose()
    {
        Dispose(true);
        GC.SuppressFinalize(this);
    }

    private void Dispose(bool disposing)
    {
        if (!_disposed)
        {
            if (disposing)
            {
                // 清理活跃备份记录
                var backupTaskKey = $"{_database}_{_today}";
                if (ActiveBackups.TryRemove(backupTaskKey, out var removedInstanceId))
                {
                    _logger.LogInformation(
                        "实例销毁,清理备份任务记录:{BackupTaskKey} -> {InstanceId}",
                        backupTaskKey,
                        removedInstanceId
                    );
                }
            }
            _disposed = true;
        }
    }

    ~BackupManager()
    {
        Dispose(false);
    }

    private BsonValue GetBsonValueFromString(string value, string type)
    {
        try
        {
            BsonValue bsonValue = type switch
            {
                "ObjectId" => new ObjectId(value),
                "String" => value,
                "Int32" => int.Parse(value, System.Globalization.CultureInfo.InvariantCulture),
                "Int64" => long.Parse(value, System.Globalization.CultureInfo.InvariantCulture),
                "Double" => double.Parse(value, System.Globalization.CultureInfo.InvariantCulture),
                "DateTime" => DateTime.Parse(
                    value,
                    System.Globalization.CultureInfo.InvariantCulture
                ),
                "Boolean" => bool.Parse(value),
                _ => value, // 默认作为字符串处理
            };
            _logger.LogDebug(
                "转换ID值:Value={Value}, Type={Type}, Result={Result}",
                value,
                type,
                bsonValue
            );
            return bsonValue;
        }
        catch (Exception e)
        {
            _logger.LogError(e, "转换ID值失败:Value={Value}, Type={Type}", value, type);
            return value; // 转换失败时返回原始字符串
        }
    }

    private static string Sha256Encrypt(string input)
    {
        var buffer = Encoding.UTF8.GetBytes(input);
        var hashBytes = SHA256.HashData(buffer);
        var key = new StringBuilder();
        foreach (var item in hashBytes)
        {
            key.Append($"{item:X2}");
        }
        return key.ToString();
    }

    /// <summary>
    /// 创建一个zip文件
    /// </summary>
    /// <param name="outPathname">zip文件路径</param>
    /// <param name="password">密码</param>
    /// <param name="folderName">要压缩的文件夹路径</param>
    private static void CreateZip(string outPathname, string password, string folderName)
    {
        using var fsOut = File.Create(outPathname);
        using var zipStream = new ZipOutputStream(fsOut);
        //0-9, 9 是最高压缩级别
        zipStream.SetLevel(9);

        // 密码
        zipStream.Password = password;

        var folderOffset = folderName.Length + (folderName.EndsWith(Path.PathSeparator) ? 0 : 1);

        CompressFolder(folderName, zipStream, folderOffset);
    }

    /// <summary>
    /// 递归压缩文件夹
    /// </summary>
    /// <param name="path"></param>
    /// <param name="zipStream"></param>
    /// <param name="folderOffset"></param>
    private static void CompressFolder(string path, ZipOutputStream zipStream, int folderOffset)
    {
        var files = Directory.GetFiles(path);

        foreach (var filename in files)
        {
            var fi = new FileInfo(filename);

            // 根据文件夹名称命名zip文件名
            var entryName = filename[folderOffset..];

            // 修正不同系统造成的不同目录的斜杠问题
            entryName = ZipEntry.CleanName(entryName);

            var newEntry = new ZipEntry(entryName)
            {
                // Zip 存储2秒的颗粒度
                DateTime = fi.LastWriteTime,
                // 指定 AESKeySize 会使用 AES 加密。
                // 允许的值为 0(关闭)、128 或 256。
                // 如果使用 AES,则需要 ZipOutputStream 上的密码
                AESKeySize = 256,
                /*
                 * 为了让WinXP和Server2003中内置的解压程序、WinZip 8、Java和其他旧代码能够解压缩这个zip文件,你需要执行以下操作之一:
                 *      指定 UseZip64.Off 或设置文件大小。
                 *      如果文件大小超过4GB,或者不需要WinXP内置的兼容性,那么不需要上述操作,但是生成的zip文件将是Zip64格式,而并非所有工具都能理解这种格式。
                 */
                // zipStream.UseZip64 = UseZip64.Off;
                Size = fi.Length,
            };

            zipStream.PutNextEntry(newEntry);

            // Zip 缓冲区
            var buffer = new byte[4096];
            using (var fsInput = File.OpenRead(filename))
            {
                StreamUtils.Copy(fsInput, zipStream, buffer);
            }

            zipStream.CloseEntry();
        }

        // 递归
        var folders = Directory.GetDirectories(path);
        foreach (var folder in folders)
        {
            CompressFolder(folder, zipStream, folderOffset);
        }
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

这是一个用于MongoDB数据库备份的管理类 BackupManager。让我详细解释一下这个代码的功能和结构:

主要功能

这个类提供了完整的MongoDB数据库备份解决方案,包括:

  • 数据导出
  • 断点续传
  • 数据压缩加密
  • 并发控制
  • 错误重试

核心特性

1. 实例管理和并发控制

private static readonly ConcurrentDictionary<string, string> ActiveBackups = new();
  • 使用静态字典跟踪活跃的备份任务
  • 防止同一数据库被多个实例同时备份
  • 每个实例都有唯一ID标识

2. 断点续传机制

private record CheckpointData(
    string CollectionName,
    string LastDocumentId, 
    string IdType,
    bool IsCompleted
);
  • 记录每个集合的备份进度
  • 支持从上次中断的位置继续备份
  • 每1000个文档更新一次检查点

3. 数据备份流程

主要备份方法:

public async Task<BackupResult> BackupAsync(Dictionary<string, string>? filter = null)

备份步骤:

  1. 获取所有集合列表
  2. 并行处理每个集合
  3. 对每个集合应用过滤条件
  4. 使用Pipeline方式高效写入BSON文件
  5. 压缩所有文件为加密ZIP

4. 错误处理和重试

  • 最大重试次数:15次
  • 重试间隔:5秒
  • 每个集合独立重试,不影响其他集合的备份

5. 高性能数据处理

var pipeWriter = PipeWriter.Create(fileStream, new StreamPipeWriterOptions(leaveOpen: true));
  • 使用System.IO.Pipelines提高I/O性能
  • 异步枚举处理大量数据
  • 内存友好的流式处理

6. 数据压缩和加密

private static string Sha256Encrypt(string input)
CreateZip(zipPath, password, _backupPath);
  • 使用SHA256生成压缩密码
  • AES-256加密压缩文件
  • 最高压缩级别(Level 9)

关键技术点

1. ID类型处理

支持多种MongoDB ID类型:

  • ObjectId
  • String
  • Int32/Int64
  • Double
  • DateTime
  • Boolean

2. 资源管理

实现了IDisposable接口,确保:

  • 清理活跃备份记录
  • 释放数据库连接
  • 防止资源泄露

3. 取消令牌支持

全程支持CancellationToken,可以优雅地取消备份操作。

使用场景

这个备份管理器特别适合:

  • 大型MongoDB数据库的定期备份
  • 需要断点续传的长时间备份任务
  • 对备份文件有加密要求的场景
  • 需要并行处理多个集合的高效备份

返回结果

return new BackupResult(password, zipPath, databaseName);

返回包含密码、文件路径和数据库名称的备份结果。

这是一个设计良好、功能完整的企业级数据库备份解决方案,考虑了性能、可靠性和安全性等多个方面。

评论加载中...