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)
备份步骤:
- 获取所有集合列表
- 并行处理每个集合
- 对每个集合应用过滤条件
- 使用Pipeline方式高效写入BSON文件
- 压缩所有文件为加密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);
返回包含密码、文件路径和数据库名称的备份结果。
这是一个设计良好、功能完整的企业级数据库备份解决方案,考虑了性能、可靠性和安全性等多个方面。
评论加载中...