diff --git a/QuestionLibraryMQConsumer/Configuration/AliyunOSSOptions.cs b/QuestionLibraryMQConsumer/Configuration/AliyunOSSOptions.cs new file mode 100644 index 0000000..aaf2664 --- /dev/null +++ b/QuestionLibraryMQConsumer/Configuration/AliyunOSSOptions.cs @@ -0,0 +1,17 @@ +namespace QuestionLibraryMQConsumer.Configuration; + +public class AliyunOSSOptions +{ + public const string SectionName = "AliyunOSSConfigs"; + + public string AccessKeyID { get; set; } = string.Empty; + public string AccessKeySecret { get; set; } = string.Empty; + public string VodBucketName { get; set; } = string.Empty; + public string BucketName { get; set; } = string.Empty; + public string Region { get; set; } = string.Empty; + public string RoleArn { get; set; } = string.Empty; + public int DurationSeconds { get; set; } + public string Endpoint { get; set; } = string.Empty; + public string ProjectName { get; set; } = string.Empty; + public string Domain { get; set; } = string.Empty; +} diff --git a/QuestionLibraryMQConsumer/Entities/Book.cs b/QuestionLibraryMQConsumer/Entities/Book.cs index b41c16d..48d5b7f 100644 --- a/QuestionLibraryMQConsumer/Entities/Book.cs +++ b/QuestionLibraryMQConsumer/Entities/Book.cs @@ -194,7 +194,7 @@ namespace QuestionLibraryMQConsumer.Entities /// Default: /// Nullable:True /// - [SugarColumn(ColumnName = "source")] - public string Source { get; set; } + [SugarColumn(ColumnName = "source_id")] + public long SourceId { get; set; } } } diff --git a/QuestionLibraryMQConsumer/Entities/OtherSystem.cs b/QuestionLibraryMQConsumer/Entities/OtherSystem.cs new file mode 100644 index 0000000..5c5f1b6 --- /dev/null +++ b/QuestionLibraryMQConsumer/Entities/OtherSystem.cs @@ -0,0 +1,50 @@ +using SqlSugar; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace QuestionLibraryMQConsumer.Entities +{ + /// + ///其他系统 + /// + [SugarTable("other_system")] + public partial class OtherSystem : SqlSugarBaseEntity + { + /// + /// Desc:主键id + /// Default: + /// Nullable:False + /// + [SugarColumn(IsPrimaryKey = true, ColumnName = "id")] + public long Id { get; set; } + + /// + /// Desc:Code + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "code")] + public string Code { get; set; } + + /// + /// Desc:请求域名 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "domain")] + public string Domain { get; set; } = string.Empty; + + + /// + /// Desc:桶名称 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "bucket_name")] + public string BucketName { get; set; } = string.Empty; + + } +} diff --git a/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs b/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs index 7f8428f..1db864c 100644 --- a/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs +++ b/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs @@ -4,6 +4,7 @@ using QuestionLibraryMQConsumer.Configuration; using QuestionLibraryMQConsumer.Consumers; using QuestionLibraryMQConsumer.Data; using QuestionLibraryMQConsumer.Handlers; +using QuestionLibraryMQConsumer.Services; namespace QuestionLibraryMQConsumer.Extensions; @@ -19,6 +20,9 @@ public static class ServiceCollectionExtensions var connectionString = configuration.GetConnectionString("zybDb") ?? throw new InvalidOperationException("Connection string 'zybDb' is not configured"); services.AddSingleton(new QuestionLibraryDb(connectionString)); + services.Configure(configuration.GetSection(AliyunOSSOptions.SectionName)); + services.AddSingleton(); + services.AddSingleton(); services.AddTransient(); diff --git a/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs b/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs index 3da4236..7bdd759 100644 --- a/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs +++ b/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs @@ -1,6 +1,7 @@ using QuestionLibraryMQConsumer.Data; using QuestionLibraryMQConsumer.Entities; using QuestionLibraryMQConsumer.Models; +using QuestionLibraryMQConsumer.Services; using RabbitMQ.Client; using System; using System.Collections.Generic; @@ -20,6 +21,7 @@ namespace QuestionLibraryMQConsumer.Handlers private readonly ILogger _logger; private readonly QuestionLibraryDb _db; + private readonly IOssService _ossService; public string RoutingKey => "questionLibrary.book.created"; public Type MessageType => typeof(BookCreatedMessage); @@ -27,16 +29,18 @@ namespace QuestionLibraryMQConsumer.Handlers public BookCreatedHandler( ILogger logger, - QuestionLibraryDb db) + QuestionLibraryDb db, + IOssService ossService) { _logger = logger; _db = db; + _ossService = ossService; } public async Task HandleJsonAsync(string messageJson, CancellationToken cancellationToken) { - _logger.LogInformation("[Handler入口] HandleJsonAsync被调用, messageJson长度: {Length}", messageJson?.Length ?? 0); - + _logger.LogInformation("[Handler入口] HandleJsonAsync被调用,messageJson长度:{msg} messageJson长度: {Length}", messageJson, messageJson?.Length ?? 0); + var message = JsonSerializer.Deserialize(messageJson, _jsonOptions); if (message == null) @@ -63,6 +67,48 @@ namespace QuestionLibraryMQConsumer.Handlers return; } + // 查询外部系统配置获取源bucket信息 + var otherSystem = await _db.Db.Queryable() + .FirstAsync(s => s.Id == message.SourceId); + + string pdfUrl = message.PdfUrl ?? string.Empty; + string coverUrl = message.CoverUrl ?? string.Empty; + string backCoverUrl = message.BackCoverUrl ?? string.Empty; + + // 如果找到外部系统配置,执行OSS跨bucket复制 + if (otherSystem != null && !string.IsNullOrEmpty(otherSystem.BucketName) && !string.IsNullOrEmpty(otherSystem.Domain)) + { + _logger.LogInformation("找到外部系统配置 SourceId: {SourceId}, Bucket: {Bucket}, Domain: {Domain}", + message.SourceId, otherSystem.BucketName, otherSystem.Domain); + + // 并行复制所有URL + var copyTasks = new List(); + + if (!string.IsNullOrEmpty(pdfUrl)) + { + copyTasks.Add(CopyUrlAsync("PdfUrl", pdfUrl, otherSystem, newUrl => pdfUrl = newUrl, cancellationToken)); + } + + if (!string.IsNullOrEmpty(coverUrl)) + { + copyTasks.Add(CopyUrlAsync("CoverUrl", coverUrl, otherSystem, newUrl => coverUrl = newUrl, cancellationToken)); + } + + if (!string.IsNullOrEmpty(backCoverUrl)) + { + copyTasks.Add(CopyUrlAsync("BackCoverUrl", backCoverUrl, otherSystem, newUrl => backCoverUrl = newUrl, cancellationToken)); + } + + if (copyTasks.Count > 0) + { + await Task.WhenAll(copyTasks); + } + } + else + { + _logger.LogWarning("未找到外部系统配置 SourceId: {SourceId},使用原始URL", message.SourceId); + } + var book = new Book { Id = message.BookId, @@ -70,12 +116,12 @@ namespace QuestionLibraryMQConsumer.Handlers Title = message.Subtitle, Width = message.PaperWidth, Height = message.PaperHeight, - Cover = message.CoverUrl, - PdfUrl = message.PdfUrl, - BackCover = message.BackCoverUrl, + Cover = coverUrl, + PdfUrl = pdfUrl, + BackCover = backCoverUrl, CreatedTime = DateTime.UtcNow, UpdatedTime = DateTime.UtcNow, - Source = message.Source, + SourceId = message.SourceId, Status = (int)BookStatusEnum.MissCatalog }; @@ -84,6 +130,21 @@ namespace QuestionLibraryMQConsumer.Handlers _logger.LogInformation("书籍 {BookId} 创建成功", message.BookId); } + private async Task CopyUrlAsync(string fieldName, string originalPath, OtherSystem otherSystem, Action setNewPath, CancellationToken cancellationToken) + { + try + { + var newPath = await _ossService.CopyToOwnBucketAsync(originalPath, otherSystem.BucketName, cancellationToken); + setNewPath(newPath); + _logger.LogInformation("{FieldName} OSS复制成功: {OriginalPath} -> {NewPath}", fieldName, originalPath, newPath); + } + catch (Exception ex) + { + _logger.LogError(ex, "{FieldName} OSS复制失败,使用原始路径: {OriginalPath}", fieldName, originalPath); + // 保持原始路径不变 + } + } + public async Task HandleWithManualAckAsync(string messageJson, IChannel channel, ulong deliveryTag, CancellationToken cancellationToken) { var message = JsonSerializer.Deserialize(messageJson, _jsonOptions); diff --git a/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs b/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs index 0af7106..ea52483 100644 --- a/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs +++ b/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs @@ -3,7 +3,7 @@ namespace QuestionLibraryMQConsumer.Models; public class BookCreatedMessage { /// - /// ��ϢId + /// 消息Id /// public string MessageId { get; set; } = string.Empty; /// @@ -15,31 +15,31 @@ public class BookCreatedMessage /// public string BookName { get; set; } = string.Empty; /// - /// ������ + /// 副标题 /// public string Subtitle { get; set; } = string.Empty; /// - /// ֽ�ſ��� + /// 纸张宽度 /// public float PaperWidth { get; set; } /// - /// ֽ�Ÿ߶� + /// 纸张高度 /// public float PaperHeight { get; set; } /// - /// ���� + /// 封面 /// public string? CoverUrl { get; set; } /// - /// ��� + /// 封底 /// public string? BackCoverUrl { get; set; } /// - /// �鼮PDF + /// 书籍PDF /// public string? PdfUrl { get; set; } /// - /// ԴʶжϻشMQϢַ + /// 来源识别判断回MQ消息地址 /// - public string? Source { get; set; } + public long SourceId { get; set; } } diff --git a/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj b/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj index e814550..5a000c0 100644 --- a/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj +++ b/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj @@ -12,5 +12,6 @@ + diff --git a/QuestionLibraryMQConsumer/Services/IOssService.cs b/QuestionLibraryMQConsumer/Services/IOssService.cs new file mode 100644 index 0000000..9e9df8f --- /dev/null +++ b/QuestionLibraryMQConsumer/Services/IOssService.cs @@ -0,0 +1,13 @@ +namespace QuestionLibraryMQConsumer.Services; + +public interface IOssService +{ + /// + /// 将文件从外部bucket复制到自己的bucket + /// + /// 源文件相对路径(如 book/807248715010117/book.pdf) + /// 源bucket名称 + /// 取消令牌 + /// 复制后的新URL + Task CopyToOwnBucketAsync(string sourcePath, string sourceBucket, CancellationToken cancellationToken = default); +} diff --git a/QuestionLibraryMQConsumer/Services/OssService.cs b/QuestionLibraryMQConsumer/Services/OssService.cs new file mode 100644 index 0000000..55a7f81 --- /dev/null +++ b/QuestionLibraryMQConsumer/Services/OssService.cs @@ -0,0 +1,51 @@ +using Aliyun.OSS; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; +using QuestionLibraryMQConsumer.Configuration; + +namespace QuestionLibraryMQConsumer.Services; + +public class OssService : IOssService +{ + private readonly ILogger _logger; + private readonly AliyunOSSOptions _ossOptions; + + public OssService( + ILogger logger, + IOptions ossOptions) + { + _logger = logger; + _ossOptions = ossOptions.Value; + } + + public async Task CopyToOwnBucketAsync(string sourcePath, string sourceBucket, CancellationToken cancellationToken = default) + { + if (string.IsNullOrEmpty(sourcePath)) + { + return sourcePath; + } + + try + { + // sourcePath 是相对路径,直接作为 object key 使用 + var objectKey = sourcePath.TrimStart('/'); + + // 创建OSS客户端 + var client = new OssClient(_ossOptions.Endpoint, _ossOptions.AccessKeyID, _ossOptions.AccessKeySecret); + + // 执行跨bucket复制 + var request = new CopyObjectRequest(sourceBucket, objectKey, _ossOptions.BucketName, objectKey); + var result = client.CopyObject(request); + + _logger.LogInformation("OSS跨bucket复制成功: {SourceBucket}/{ObjectKey} -> {TargetBucket}/{TargetKey}, ETag: {ETag}", + sourceBucket, objectKey, _ossOptions.BucketName, objectKey, result.ETag); + + return objectKey; + } + catch (Exception ex) + { + _logger.LogError(ex, "OSS跨bucket复制失败: {SourcePath}", sourcePath); + return sourcePath; // 失败时返回原始路径 + } + } +} diff --git a/QuestionLibraryMQConsumer/appsettings.json b/QuestionLibraryMQConsumer/appsettings.json index dbd5736..d9dbdc3 100644 --- a/QuestionLibraryMQConsumer/appsettings.json +++ b/QuestionLibraryMQConsumer/appsettings.json @@ -52,5 +52,17 @@ } } ] + }, + "AliyunOSSConfigs": { + "AccessKeyID": "LTAI5tEBXGewpHSLiSxyx6Bf", + "AccessKeySecret": "w29b8wkw6XQVL8GWXgp3ZesgYeDKvf", + "VodBucketName": "outin-5277bbb52bec11f08dbd00163e169e2b.oss-cn-beijing.aliyuncs.com", + "BucketName": "qyzh-qb", + "Region": "beijing", + "RoleArn": "acs:ram::1064745380176636:role/aliyunosstokengeneratorrole", + "DurationSeconds": 3600, //过期时间(秒) + "Endpoint": "oss-cn-beijing.aliyuncs.com", + "ProjectName": "scs", + "Domain": "https://qyzh-qb.oss-cn-beijing.aliyuncs.com/" } }