feat: 新增COS媒体文件迁移全流程功能,优化实时上传逻辑

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,新增迁移范围说明
This commit is contained in:
jqb 2026-06-26 09:50:39 +08:00
parent 1730d7c746
commit 649f71edc9
10 changed files with 272 additions and 36 deletions

View File

@ -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 的构建、上传、部署和自检。
---
## 服务器信息

View File

@ -137,7 +137,8 @@ public class ArchiveCallbackController {
try {
long lastSeq = archivePullService.getLastSeq();
List<ArchiveMessage> 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);

View File

@ -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<Integer> migrate(@RequestParam(defaultValue = "50") Integer batchSize) {
public Result<Integer> 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<Map<String, Object>> 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<String, Object> 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<Map<String, Object>> migrateStatus(
@RequestParam(defaultValue = "7") Integer recentDays) {
Long minMsgtime = calcMinMsgtime(recentDays);
Map<String, Object> 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);

View File

@ -113,6 +113,19 @@ public interface ArchiveMessageMapper extends BaseMapper<ArchiveMessage> {
"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<ArchiveMessage> 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<ArchiveMessage> 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);
}

View File

@ -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<ArchiveMessage> messages) {
log.info("开始保存消息, 共{}条", messages.size());
saveAndNotify(messages, false);
}
/**
* 保存消息并可选立即预上传到 COS
*
* @param messages 消息列表
* @param immediateCosUpload 是否对图片/语音立即预上传到 COS(回调实时拉取时建议开启)
*/
public void saveAndNotify(List<ArchiveMessage> 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() {

View File

@ -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,11 +300,11 @@ public class CosStorageService {
}
int count = 0;
int offset = 0;
while (true) {
List<ArchiveMessage> messages = archiveMessageMapper.selectMediaMessagesWithoutCos(offset, batchSize);
// 从指定偏移开始迁移,id 倒序(最新的在前)
List<ArchiveMessage> messages = archiveMessageMapper.selectMediaMessagesWithoutCos(offset, batchSize, minMsgtime);
if (messages == null || messages.isEmpty()) {
break;
log.info("没有需要迁移的媒体文件,offset={}", offset);
return 0;
}
for (ArchiveMessage message : messages) {
try {
@ -280,14 +316,7 @@ public class CosStorageService {
log.error("迁移媒体文件失败: id={}", message.getId(), e);
}
}
offset += batchSize;
// 简单限流,避免触发接口频率限制
try {
Thread.sleep(50);
} catch (InterruptedException ignored) {
}
}
log.info("批量迁移完成,共迁移 {} 个文件", count);
log.info("本批迁移完成,offset={}, 成功 {} / {} 个文件", offset, count, messages.size());
return count;
}

32
deploy-archive-service.sh Normal file
View File

@ -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 "发布完成。"

View File

@ -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',

View File

@ -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() {
<Card style={{ marginTop: 16 }} title="历史媒体文件迁移">
<p style={{ color: '#666', marginBottom: 16 }}>
将已下载到本地的会话存档媒体文件批量上传到腾讯云 COS。上传完成后,后续访问将直接走 COS 链接。
<br />
当前仅迁移图片和语音,视频/文件暂不迁移。每次处理 10 条,可多次点击。
</p>
<Button
type="primary"

74
migrate-all-to-cos.sh Normal file
View File

@ -0,0 +1,74 @@
#!/bin/bash
# 后台迁移最近 N 天的图片/语音到腾讯云 COS
# 用法:在服务器上执行 nohup ./migrate-all-to-cos.sh > /tmp/migrate-cos.log 2>&1 &
set -u
BATCH_ENDPOINT="http://localhost:8082/api/v1/archive/cos-config/migrate-batch"
STATUS_ENDPOINT="http://localhost:8082/api/v1/archive/cos-config/migrate-status"
BATCH_SIZE=10
RECENT_DAYS=3
SLEEP_SECONDS=0
LOG_EVERY=50
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 开始迁移最近 ${RECENT_DAYS} 天的图片/语音到 COS,批次大小 ${BATCH_SIZE}"
# 先获取当前未迁移总数(最近 N 天)
status=$(curl -s --max-time 30 "${STATUS_ENDPOINT}?recentDays=${RECENT_DAYS}")
unmapped=$(echo "${status}" | grep -oP '"unmapped":\s*\K[0-9]+' || echo "0")
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 当前未迁移图片/语音(最近 ${RECENT_DAYS} 天):${unmapped}"
if [ "${unmapped}" -eq 0 ]; then
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 没有需要迁移的数据,退出"
exit 0
fi
# 计算需要扫描的最大偏移量
maxOffset=$(( (unmapped / BATCH_SIZE + 5) * BATCH_SIZE ))
pass=0
totalMigrated=0
while true; do
pass=$((pass + 1))
passMigrated=0
round=0
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 开始第 ${pass} 轮扫描(最近 ${RECENT_DAYS} 天),maxOffset=${maxOffset}"
for offset in $(seq 0 ${BATCH_SIZE} ${maxOffset}); do
round=$((round + 1))
result=$(curl -s -X POST --max-time 60 "${BATCH_ENDPOINT}?offset=${offset}&batchSize=${BATCH_SIZE}&recentDays=${RECENT_DAYS}")
code=$(echo "${result}" | grep -oP '"code":\s*\K[0-9]+' || echo "")
count=$(echo "${result}" | grep -oP '"count":\s*\K[0-9]+' || echo "")
if [ "${code}" != "0" ] || [ -z "${count}" ]; then
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 第 ${pass} 轮第 ${round} 批(offset=${offset})调用失败: ${result}"
else
passMigrated=$((passMigrated + count))
totalMigrated=$((totalMigrated + count))
if [ "${count}" -gt 0 ]; then
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 第 ${pass} 轮 offset=${offset} 迁移 ${count} 个,本轮累计 ${passMigrated},总累计 ${totalMigrated}"
fi
fi
# 每 LOG_EVERY 批输出一次进度,避免日志看起来卡住
if [ $((round % LOG_EVERY)) -eq 0 ]; then
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 第 ${pass} 轮进度:已处理 ${round} 批,当前 offset=${offset},本轮迁移 ${passMigrated} 个,总累计 ${totalMigrated} 个"
fi
sleep "${SLEEP_SECONDS}"
done
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 第 ${pass} 轮扫描结束,本轮迁移 ${passMigrated} 个,总累计 ${totalMigrated} 个"
# 本轮没有成功,结束任务
if [ "${passMigrated}" -eq 0 ]; then
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 本轮无新增迁移,任务结束"
break
fi
done
# 最终状态
status=$(curl -s --max-time 30 "${STATUS_ENDPOINT}?recentDays=${RECENT_DAYS}")
mapped=$(echo "${status}" | grep -oP '"mapped":\s*\K[0-9]+' || echo "0")
unmapped=$(echo "${status}" | grep -oP '"unmapped":\s*\K[0-9]+' || echo "0")
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 最终状态:已映射 ${mapped},最近 ${RECENT_DAYS} 天未映射 ${unmapped}"