From 649f71edc9cc564e299fc1d23a3ab3acc9cd0260 Mon Sep 17 00:00:00 2001 From: jqb Date: Fri, 26 Jun 2026 09:50:39 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=96=B0=E5=A2=9ECOS=E5=AA=92=E4=BD=93?= =?UTF-8?q?=E6=96=87=E4=BB=B6=E8=BF=81=E7=A7=BB=E5=85=A8=E6=B5=81=E7=A8=8B?= =?UTF-8?q?=E5=8A=9F=E8=83=BD=EF=BC=8C=E4=BC=98=E5=8C=96=E5=AE=9E=E6=97=B6?= =?UTF-8?q?=E4=B8=8A=E4=BC=A0=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 调整archive_media_files表sdkfileid字段为TEXT类型适配长内容 2. 新增一键部署脚本deploy-archive-service.sh和全量迁移脚本migrate-all-to-cos.sh 3. 重构COS迁移接口,支持按批次、按时间范围迁移图片和语音文件 4. 新增迁移进度查询接口,支持统计已迁移/未迁移数量 5. 优化实时回调逻辑,新增即时预上传COS功能 6. 更新AGENTS.md文档,补充部署说明和脚本使用说明 7. 调整前端迁移批次大小为10,新增迁移范围说明 --- AGENTS.md | 9 ++- .../controller/ArchiveCallbackController.java | 3 +- .../controller/CosConfigController.java | 57 +++++++++++++- .../archive/mapper/ArchiveMessageMapper.java | 17 ++++- .../archive/service/ArchivePullService.java | 35 ++++++++- .../archive/service/CosStorageService.java | 75 +++++++++++++------ deploy-archive-service.sh | 32 ++++++++ deploy_cos_tables.sql | 2 +- frontend/admin/src/pages/CosConfig.tsx | 4 +- migrate-all-to-cos.sh | 74 ++++++++++++++++++ 10 files changed, 272 insertions(+), 36 deletions(-) create mode 100644 deploy-archive-service.sh create mode 100644 migrate-all-to-cos.sh diff --git a/AGENTS.md b/AGENTS.md index a40cd63..de4e436 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -110,16 +110,21 @@ cd backend/auth-service && mvn clean package -DskipTests # 2. 将 JAR 包同步到服务器的 deploy-package 对应目录 scp -o StrictHostKeyChecking=no -i root.pem target/*.jar root@8.133.162.25:/www/wwwroot/deploy-package/backend/auth-service/ -# 3. SSH 登录服务器,重启对应容器 +# 3. SSH 登录服务器,重新构建镜像并启动容器 # ⚠️ 注意:Dockerfile 中 COPY 的是 app.jar,必须先把新 jar 复制为 app.jar +# ⚠️ docker-compose up 没有 --no-cache 参数,如需无缓存构建应分两步执行 ssh -o StrictHostKeyChecking=no -i root.pem root@8.133.162.25 << 'EOF' cd /www/wwwroot/deploy-package/backend/auth-service cp auth-service-*.jar app.jar cd /www/wwwroot/deploy-package -docker-compose up -d --build --no-cache auth-service +docker-compose build --no-cache auth-service +docker-compose up -d auth-service EOF ``` +### 一键脚本 +项目根目录已提供 `deploy-archive-service.sh`,可自动完成 archive-service 的构建、上传、部署和自检。 + --- ## 服务器信息 diff --git a/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveCallbackController.java b/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveCallbackController.java index 6b72145..9cca2d4 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveCallbackController.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveCallbackController.java @@ -137,7 +137,8 @@ public class ArchiveCallbackController { try { long lastSeq = archivePullService.getLastSeq(); List messages = archivePullService.pullMessages(lastSeq, 1000); - archivePullService.saveAndNotify(messages); + // 回调触发的是实时新消息,立即预上传到 COS + archivePullService.saveAndNotify(messages, true); log.info("回调触发拉取完成,共{}条消息", messages.size()); } catch (Exception e) { log.error("回调触发拉取失败: {}", e.getMessage(), e); diff --git a/backend/archive-service/src/main/java/com/artedu/archive/controller/CosConfigController.java b/backend/archive-service/src/main/java/com/artedu/archive/controller/CosConfigController.java index 3e03c27..6e8375f 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/controller/CosConfigController.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/controller/CosConfigController.java @@ -12,7 +12,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.*; import java.time.LocalDateTime; +import java.util.HashMap; import java.util.List; +import java.util.Map; /** * 腾讯云 COS 配置管理接口 @@ -135,17 +137,64 @@ public class CosConfigController { } /** - * 批量迁移历史媒体文件到 COS + * 批量迁移历史媒体文件到 COS(前端按钮使用,始终从最新未迁移开始) */ @PostMapping("/migrate") - public Result migrate(@RequestParam(defaultValue = "50") Integer batchSize) { + public Result migrate( + @RequestParam(defaultValue = "10") Integer batchSize, + @RequestParam(defaultValue = "7") Integer recentDays) { if (batchSize == null || batchSize <= 0 || batchSize > 200) { - batchSize = 50; + batchSize = 10; } - int count = cosStorageService.migrateLocalFiles(batchSize); + Long minMsgtime = calcMinMsgtime(recentDays); + int count = cosStorageService.migrateLocalFiles(0, batchSize, minMsgtime); return Result.success(count); } + /** + * 批量迁移指定偏移的一批(后台全量扫描脚本使用) + */ + @PostMapping("/migrate-batch") + public Result> migrateBatch( + @RequestParam(defaultValue = "0") Integer offset, + @RequestParam(defaultValue = "10") Integer batchSize, + @RequestParam(defaultValue = "7") Integer recentDays) { + if (offset == null || offset < 0) { + offset = 0; + } + if (batchSize == null || batchSize <= 0 || batchSize > 200) { + batchSize = 10; + } + Long minMsgtime = calcMinMsgtime(recentDays); + int count = cosStorageService.migrateLocalFiles(offset, batchSize, minMsgtime); + Map result = new HashMap<>(); + result.put("offset", offset); + result.put("batchSize", batchSize); + result.put("count", count); + return Result.success(result); + } + + /** + * 查询 COS 迁移进度 + */ + @GetMapping("/migrate-status") + public Result> migrateStatus( + @RequestParam(defaultValue = "7") Integer recentDays) { + Long minMsgtime = calcMinMsgtime(recentDays); + Map result = new HashMap<>(); + result.put("mapped", cosStorageService.countMappedMediaFiles()); + result.put("unmapped", cosStorageService.countUnmappedMediaFiles(minMsgtime)); + result.put("recentDays", recentDays); + return Result.success(result); + } + + private Long calcMinMsgtime(Integer recentDays) { + if (recentDays == null || recentDays <= 0) { + return null; + } + return System.currentTimeMillis() - (long) recentDays * 24 * 60 * 60 * 1000; + } + private CosConfigVO toVO(CosConfig config, boolean maskSecretKey) { CosConfigVO vo = new CosConfigVO(); BeanUtils.copyProperties(config, vo); diff --git a/backend/archive-service/src/main/java/com/artedu/archive/mapper/ArchiveMessageMapper.java b/backend/archive-service/src/main/java/com/artedu/archive/mapper/ArchiveMessageMapper.java index e174a30..eed85ee 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/mapper/ArchiveMessageMapper.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/mapper/ArchiveMessageMapper.java @@ -113,6 +113,19 @@ public interface ArchiveMessageMapper extends BaseMapper { "LEFT JOIN archive_media_files f ON m.id = f.archive_message_id AND f.status = 1 " + "WHERE m.media_data IS NOT NULL AND m.media_data != '' " + "AND f.id IS NULL " + - "ORDER BY m.id LIMIT #{offset}, #{limit}") - List selectMediaMessagesWithoutCos(@Param("offset") int offset, @Param("limit") int limit); + "AND m.msgtype IN ('image', 'voice') " + + "AND (#{minMsgtime} IS NULL OR m.msgtime >= #{minMsgtime}) " + + "ORDER BY m.id DESC LIMIT #{offset}, #{limit}") + List selectMediaMessagesWithoutCos(@Param("offset") int offset, @Param("limit") int limit, @Param("minMsgtime") Long minMsgtime); + + /** + * 尚未上传到 COS 的图片/语音总数 + */ + @Select("SELECT COUNT(*) FROM archive_messages m " + + "LEFT JOIN archive_media_files f ON m.id = f.archive_message_id AND f.status = 1 " + + "WHERE m.media_data IS NOT NULL AND m.media_data != '' " + + "AND f.id IS NULL " + + "AND m.msgtype IN ('image', 'voice') " + + "AND (#{minMsgtime} IS NULL OR m.msgtime >= #{minMsgtime})") + Long selectUnmappedMediaCount(@Param("minMsgtime") Long minMsgtime); } diff --git a/backend/archive-service/src/main/java/com/artedu/archive/service/ArchivePullService.java b/backend/archive-service/src/main/java/com/artedu/archive/service/ArchivePullService.java index d13dc93..53a39a8 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/service/ArchivePullService.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/service/ArchivePullService.java @@ -9,6 +9,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Lazy; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Service; @@ -67,6 +68,10 @@ public class ArchivePullService { @Autowired private ArchiveRiskService archiveRiskService; + @Autowired + @Lazy + private CosStorageService cosStorageService; + private static final String LAST_SEQ_KEY = "archive:last_seq:"; @PostConstruct @@ -466,8 +471,19 @@ public class ArchivePullService { } public void saveAndNotify(List messages) { - log.info("开始保存消息, 共{}条", messages.size()); + saveAndNotify(messages, false); + } + + /** + * 保存消息并可选立即预上传到 COS + * + * @param messages 消息列表 + * @param immediateCosUpload 是否对图片/语音立即预上传到 COS(回调实时拉取时建议开启) + */ + public void saveAndNotify(List messages, boolean immediateCosUpload) { + log.info("开始保存消息, 共{}条, immediateCosUpload={}", messages.size(), immediateCosUpload); int saved = 0; + int cosUploaded = 0; for (ArchiveMessage message : messages) { try { int count = archiveMessageMapper.countByMsgId(message.getMsgid(), message.getCorpId()); @@ -494,6 +510,21 @@ public class ArchivePullService { // 风控检测(敏感词扫描等) archiveRiskService.processMessageRisk(message); + // 实时拉取时,对图片/语音立即预上传到 COS + if (immediateCosUpload && cosStorageService != null + && message.getMediaData() != null && !message.getMediaData().isEmpty() + && ("image".equals(message.getMsgtype()) || "voice".equals(message.getMsgtype()))) { + try { + String cosUrl = cosStorageService.ensureUploaded(message); + if (cosUrl != null && !cosUrl.isEmpty()) { + cosUploaded++; + log.info("消息预上传到 COS 成功: msgid={}, url={}", message.getMsgid(), cosUrl); + } + } catch (Exception ex) { + log.warn("消息预上传到 COS 失败: msgid={}, error={}", message.getMsgid(), ex.getMessage()); + } + } + rabbitTemplate.convertAndSend( RabbitConfig.ARCHIVE_EXCHANGE, RabbitConfig.ROUTING_KEY_NEW_MESSAGE, @@ -510,7 +541,7 @@ public class ArchivePullService { log.error("保存消息失败: msgid={}, error={}", message.getMsgid(), e.getMessage(), e); } } - log.info("保存消息完成, 成功{}条/共{}条", saved, messages.size()); + log.info("保存消息完成, 成功{}条/共{}条, 预上传COS{}条", saved, messages.size(), cosUploaded); } public long getLastSeq() { diff --git a/backend/archive-service/src/main/java/com/artedu/archive/service/CosStorageService.java b/backend/archive-service/src/main/java/com/artedu/archive/service/CosStorageService.java index 1e1c2df..05061f5 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/service/CosStorageService.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/service/CosStorageService.java @@ -198,6 +198,13 @@ public class CosStorageService { return existingUrl; } + // 目前只把图片和语音上传到 COS,视频/文件等大媒体走本地下载 + String msgtype = message.getMsgtype(); + if (!"image".equals(msgtype) && !"voice".equals(msgtype)) { + log.debug("非图片/语音媒体,跳过 COS 上传: id={}, msgtype={}", message.getId(), msgtype); + return null; + } + String mediaData = message.getMediaData(); if (mediaData == null || mediaData.isEmpty()) { return null; @@ -254,9 +261,38 @@ public class CosStorageService { } /** - * 批量迁移本地历史媒体文件到 COS + * 已映射到 COS 的媒体文件数量 */ - public int migrateLocalFiles(int batchSize) { + public long countMappedMediaFiles() { + try { + return archiveMediaFileMapper.selectCount(null); + } catch (Exception e) { + log.error("查询已映射数量失败", e); + return 0; + } + } + + /** + * 尚未映射到 COS 的图片/语音数量 + */ + public long countUnmappedMediaFiles(Long minMsgtime) { + try { + return archiveMessageMapper.selectUnmappedMediaCount(minMsgtime); + } catch (Exception e) { + log.error("查询未映射数量失败", e); + return 0; + } + } + + /** + * 批量迁移本地历史媒体文件到 COS(每次调用只处理一批,避免请求超时) + * + * @param offset 分页偏移 + * @param batchSize 每批条数 + * @param minMsgtime 最早消息时间(毫秒时间戳),只迁移该时间之后的消息 + * @return 本批成功上传的条数 + */ + public int migrateLocalFiles(int offset, int batchSize, Long minMsgtime) { CosConfig config = getActiveConfig(); if (config == null || config.getEnabled() == null || config.getEnabled() != 1) { log.warn("COS 未启用,无法迁移"); @@ -264,30 +300,23 @@ public class CosStorageService { } int count = 0; - int offset = 0; - while (true) { - List messages = archiveMessageMapper.selectMediaMessagesWithoutCos(offset, batchSize); - if (messages == null || messages.isEmpty()) { - break; - } - for (ArchiveMessage message : messages) { - try { - String url = ensureUploaded(message); - if (url != null && !url.isEmpty()) { - count++; - } - } catch (Exception e) { - log.error("迁移媒体文件失败: id={}", message.getId(), e); - } - } - offset += batchSize; - // 简单限流,避免触发接口频率限制 + // 从指定偏移开始迁移,id 倒序(最新的在前) + List messages = archiveMessageMapper.selectMediaMessagesWithoutCos(offset, batchSize, minMsgtime); + if (messages == null || messages.isEmpty()) { + log.info("没有需要迁移的媒体文件,offset={}", offset); + return 0; + } + for (ArchiveMessage message : messages) { try { - Thread.sleep(50); - } catch (InterruptedException ignored) { + String url = ensureUploaded(message); + if (url != null && !url.isEmpty()) { + count++; + } + } catch (Exception e) { + log.error("迁移媒体文件失败: id={}", message.getId(), e); } } - log.info("批量迁移完成,共迁移 {} 个文件", count); + log.info("本批迁移完成,offset={}, 成功 {} / {} 个文件", offset, count, messages.size()); return count; } diff --git a/deploy-archive-service.sh b/deploy-archive-service.sh new file mode 100644 index 0000000..2676bab --- /dev/null +++ b/deploy-archive-service.sh @@ -0,0 +1,32 @@ +#!/bin/bash +# archive-service 手动发布脚本 +# 用法:在项目根目录执行 ./deploy-archive-service.sh +set -e + +SERVER_IP="8.133.162.25" +SSH_KEY="root.pem" +REMOTE_PATH="/www/wwwroot/deploy-package/backend/archive-service" +LOCAL_JAR="backend/archive-service/target/archive-service-1.0.0-SNAPSHOT.jar" + +echo "[1/4] 本地构建 archive-service..." +cd backend/archive-service +mvn clean package -DskipTests +cd ../.. + +echo "[2/4] 上传 JAR 到服务器..." +scp -o StrictHostKeyChecking=no -i "$SSH_KEY" "$LOCAL_JAR" "root@$SERVER_IP:$REMOTE_PATH/" + +echo "[3/4] 服务器端重新构建镜像并启动容器..." +ssh -o StrictHostKeyChecking=no -i "$SSH_KEY" "root@$SERVER_IP" << EOF + cd "$REMOTE_PATH" + cp archive-service-1.0.0-SNAPSHOT.jar app.jar + cd /www/wwwroot/deploy-package + docker-compose build --no-cache archive-service + docker-compose up -d archive-service +EOF + +echo "[4/4] 等待服务就绪并自检..." +sleep 20 +curl -s "https://ai.9artedu.com/api/v1/archive/cos-config" | head -c 300 +echo "" +echo "发布完成。" diff --git a/deploy_cos_tables.sql b/deploy_cos_tables.sql index 83760a7..dc53a83 100644 --- a/deploy_cos_tables.sql +++ b/deploy_cos_tables.sql @@ -16,7 +16,7 @@ CREATE TABLE IF NOT EXISTS `archive_media_files` ( `id` BIGINT(20) PRIMARY KEY AUTO_INCREMENT COMMENT '自增主键', `msg_id` VARCHAR(255) NOT NULL COMMENT '企微消息 msgid', `archive_message_id` BIGINT(20) NOT NULL COMMENT 'archive_messages.id', - `sdkfileid` VARCHAR(512) DEFAULT NULL COMMENT '企微媒体文件 sdkfileid', + `sdkfileid` TEXT DEFAULT NULL COMMENT '企微媒体文件 sdkfileid', `cos_url` VARCHAR(1024) NOT NULL COMMENT 'COS 文件访问 URL', `file_size` BIGINT(20) DEFAULT 0 COMMENT '文件大小(字节)', `file_type` VARCHAR(32) DEFAULT NULL COMMENT '文件类型:image/voice/video/file', diff --git a/frontend/admin/src/pages/CosConfig.tsx b/frontend/admin/src/pages/CosConfig.tsx index 9967496..cb347e6 100644 --- a/frontend/admin/src/pages/CosConfig.tsx +++ b/frontend/admin/src/pages/CosConfig.tsx @@ -105,7 +105,7 @@ export default function CosConfig() { setMigrateLoading(true) setMigratedCount(null) try { - const res: any = await request.post('/v1/archive/cos-config/migrate?batchSize=50') + const res: any = await request.post('/v1/archive/cos-config/migrate?batchSize=10') setMigratedCount(res.data || 0) message.success(`批量迁移完成,共迁移 ${res.data || 0} 个文件`) } catch (e: any) { @@ -195,6 +195,8 @@ export default function CosConfig() {

将已下载到本地的会话存档媒体文件批量上传到腾讯云 COS。上传完成后,后续访问将直接走 COS 链接。 +
+ 当前仅迁移图片和语音,视频/文件暂不迁移。每次处理 10 条,可多次点击。