using Dpz.Core.Backup;
using Dpz.Core.Entity.Base;
using Dpz.Core.Infrastructure;
using Dpz.Core.Public.Entity;
using Dpz.Core.Public.ViewModel;
using Dpz.Core.Service.RepositoryService;
using Hangfire;
using JetBrains.Annotations;
using MongoDB.Bson.Serialization;
using MongoDB.Driver;
namespace Dpz.Core.Web.Jobs.Hangfire;
[UsedImplicitly]
public class BackupActivator(
IBackupService backupService,
IBackupRecordService backupRecordService,
IConfiguration configuration,
//ISafeFileService safeFileService,
ILogger<BackupActivator> logger
) : JobActivator
{
[ProlongExpirationTime]
public async Task BackupAsync()
{
var filter = new Dictionary<string, string>
{
{
nameof(MessageOutboxRecord),
SerializerFilterDocument(
new ExpressionFilterDefinition<MessageOutboxRecord>(x =>
x.Status != OutboxMessageStatus.Consumed
)
)
},
};
try
{
await BackupMainAsync("mongodb", filter);
await BackupMainAsync("AgileConfig");
}
catch (Exception e)
{
logger.LogError(e, "backup fail");
}
}
private static string SerializerFilterDocument<T>(FilterDefinition<T> filter)
where T : IBaseEntity
{
var serializerRegistry = BsonSerializer.SerializerRegistry;
var documentSerializer = serializerRegistry.GetSerializer<T>();
var serializerArgs = new RenderArgs<T>(documentSerializer, serializerRegistry);
var result = filter.Render(serializerArgs);
return result.ToString();
}
private async Task BackupMainAsync(
string connectionStringName,
Dictionary<string, string>? filter = null
)
{
logger.LogInformation("Start backing up database: {Database}", connectionStringName);
var connectionString = configuration.GetConnectionString(connectionStringName);
backupService.UploadAsync = UploadAsync;
var backupRecord = await backupService.BackupAsync(connectionString, filter);
await backupRecordService.AddRecordAsync(backupRecord);
logger.LogInformation("Database: {Database} backup complete", connectionStringName);
}
private async Task<VmBackupRecord?> UploadAsync(BackupResult result)
{
if (string.IsNullOrEmpty(result.BackupPath))
{
throw new BusinessException("backup path is null");
}
if (!File.Exists(result.BackupPath))
{
throw new BusinessException("backup file not exist");
}
var backupRootPath = configuration.GetValue<string>("BackupRootPath");
if (string.IsNullOrEmpty(backupRootPath) || !Directory.Exists(backupRootPath))
{
throw new BusinessException("backup root path is null");
}
var date = DateTime.Now;
var backupPath = new List<string> { "db", date.Year.ToString(), date.Month.ToString() };
var fileName = Path.GetFileName(result.BackupPath);
#if DEBUG
backupPath.Insert(1, "Test");
#endif
// var cloudFile = new CloudFile
// {
// PathToFile = pathToFile,
// Stream = new FileStream(
// result.BackupPath,
// FileMode.Open,
// FileAccess.Read,
// FileShare.Read,
// 1 << 12
// ),
// };
// var address = await safeFileService.UploadFileForFtpAsync(cloudFile);
// if (address == null)
// {
// throw new BusinessException("upload backup file fail");
// }
var backupPathDir = Path.Combine(backupRootPath, Path.Combine(backupPath.ToArray()));
var destination = Path.Combine(backupPathDir, fileName);
try
{
if (!Directory.Exists(backupPathDir))
{
Directory.CreateDirectory(backupPathDir);
}
destination = Path.Combine(backupPathDir, fileName);
File.Move(result.BackupPath, destination);
await using var writer = File.CreateText(destination + "__password.txt");
await writer.WriteAsync(result.ZipPassword);
await writer.FlushAsync();
writer.Close();
}
catch (Exception e)
{
logger.LogError(
e,
"move backup file fail, source: {Source}, destination: {Destination}",
result.BackupPath,
destination
);
throw;
}
return new VmBackupRecord
{
Filename = destination,
BackupTime = DateTime.Now,
Database = result.Database,
BackupPassword = result.ZipPassword,
};
}
}
⚠⚠ 以下内容为AI分析的结果,请根据实际情况进行判断。
代码解释
这是一个用于数据库备份的 Hangfire 后台任务激活器类,主要功能是定期备份数据库并将备份文件保存到指定位置。
类结构
主要依赖注入
IBackupService: 备份服务IBackupRecordService: 备份记录服务IConfiguration: 配置服务ILogger: 日志记录器
核心方法
1. BackupAsync() - 主备份方法
[ProlongExpirationTime]
public async Task BackupAsync()
- 使用
[ProlongExpirationTime]特性延长 Hangfire 任务过期时间 - 定义过滤器字典,针对
MessageOutboxRecord表只备份状态不为Consumed的记录 - 依次备份两个数据库:
- "mongodb" - 应用过滤器
- "AgileConfig" - 不应用过滤器
- 异常统一捕获并记录日志
2. SerializerFilterDocument<T>() - 过滤器序列化
private static string SerializerFilterDocument<T>(FilterDefinition<T> filter)
- 将 MongoDB 的
FilterDefinition对象序列化为字符串 - 用于将过滤条件持久化或传递
3. BackupMainAsync() - 执行备份逻辑
private async Task BackupMainAsync(string connectionStringName, Dictionary<string, string>? filter = null)
- 从配置中获取连接字符串
- 设置备份服务的上传回调方法
- 执行备份并保存备份记录
- 记录开始和完成日志
4. UploadAsync() - 处理备份文件
private async Task<VmBackupRecord?> UploadAsync(BackupResult result)
主要逻辑:
验证备份文件
- 检查备份路径是否为空
- 检查备份文件是否存在
- 验证备份根目录配置
构建目标路径
{BackupRootPath}/db/{Year}/{Month}/{FileName}- DEBUG 模式下会插入 "Test" 子目录
文件操作
- 创建目标目录(如不存在)
- 移动备份文件到目标位置
- 创建密码文件(
{FileName}__password.txt)保存压缩密码
返回备份记录
- 返回包含文件名、备份时间、数据库名、密码的
VmBackupRecord对象
- 返回包含文件名、备份时间、数据库名、密码的
注释掉的代码
代码中包含被注释的 FTP 上传逻辑,说明可能之前使用云存储,现在改为本地文件系统存储。
关键特性
- ✅ 支持条件过滤备份(可选择性备份某些表的部分数据)
- ✅ 自动按年月组织备份文件
- ✅ 备份文件带密码保护,密码单独存储
- ✅ 完整的异常处理和日志记录
- ✅ 使用 Hangfire 进行后台任务调度
使用场景
适用于需要定期自动备份 MongoDB 和 AgileConfig 配置数据库的应用系统。
AI 正在分析代码…
评论加载中...