diff --git a/src/main/java/com/imutil/controller/AdminController.java b/src/main/java/com/imutil/controller/AdminController.java index d951b03..f08281c 100644 --- a/src/main/java/com/imutil/controller/AdminController.java +++ b/src/main/java/com/imutil/controller/AdminController.java @@ -634,6 +634,32 @@ public class AdminController { ctx.redirect(basePath + "/admin/sync?msg=smdone"); } + /** 同步单聊消息:需先同步用户(群成员) */ + @Post + @Mapping("/sync/c2cmsg") + public void syncC2CMsg(@Param(required = false) Long sourceAppId, + @Param(defaultValue = "") String tenantId, + Context ctx) throws Throwable { + if (sourceAppId == null && (tenantId == null || tenantId.isEmpty())) { + ctx.redirect(basePath + "/admin/sync?msg=syncopt"); + return; + } + syncService.syncC2CMessages(sourceAppId, tenantId); + ctx.redirect(basePath + "/admin/sync?msg=scdone"); + } + + /** 补拉对账:重跑群+单聊(幂等补缺失),仅主应用 + 指定租户 */ + @Post + @Mapping("/sync/check") + public void syncCheck(@Param String tenantId, Context ctx) throws Throwable { + if (tenantId == null || tenantId.isEmpty()) { + ctx.redirect(basePath + "/admin/sync?msg=syncopt"); + return; + } + syncService.checkAndPull(tenantId); + ctx.redirect(basePath + "/admin/sync?msg=chkdone"); + } + // ==================== 公共:构造页面模型 ==================== /** diff --git a/src/main/java/com/imutil/service/DispatchService.java b/src/main/java/com/imutil/service/DispatchService.java index cd3b96e..a843ce6 100644 --- a/src/main/java/com/imutil/service/DispatchService.java +++ b/src/main/java/com/imutil/service/DispatchService.java @@ -1,7 +1,9 @@ package com.imutil.service; import com.imutil.entity.DistQueue; +import com.imutil.entity.ImMessage; +import java.time.OffsetDateTime; import java.util.List; /** @@ -32,4 +34,24 @@ public interface DispatchService { * @return 重置条数 */ int recoverStuck(); + + /** + * 写入分发队列(默认 nextRetryAt=now,立即可被消费分发) + *
+ * 取租户 callbackUrl 作为目标;租户未配置 callbackUrl 则跳过(不报错,本地留底仍有效)。 + * + * @param tenantId 归属租户 + * @param msg 已落库消息(payload 取其 msgKey/convType/convId/from/to/msgTime/msgType/msgBody/source) + */ + void enqueue(String tenantId, ImMessage msg); + + /** + * 写入分发队列并指定下次分发时间(限速平滑分发用) + *
+ * 同步批量入队时按序号错开 nextRetryAt(delay = seq / dispatchQps 秒), + * DispatchWorker 按时间消费,避免历史消息瞬时打爆业务系统 callbackUrl。 + * + * @param nextRetryAt 下次可分发时间 + */ + void enqueue(String tenantId, ImMessage msg, OffsetDateTime nextRetryAt); } diff --git a/src/main/java/com/imutil/service/SyncService.java b/src/main/java/com/imutil/service/SyncService.java index ee7fb64..61e5873 100644 --- a/src/main/java/com/imutil/service/SyncService.java +++ b/src/main/java/com/imutil/service/SyncService.java @@ -9,8 +9,9 @@ import com.imutil.entity.MigrateTask; * 本服务定位为只读拉取:把指定腾讯 IM 应用里的群组/群成员/群消息拉到本地表, * 供管理后台按租户查看,不在腾讯侧产生任何写操作。 *
- * 腾讯 IM 现实约束:无"全量用户列表"接口(用户靠群成员反推);无"全量 C2C 会话列表"接口 - * (单聊历史消息本轮不做全量,回调增量已由 {@link com.imutil.service.PullService} 覆盖)。 + * 腾讯 IM 现实约束:无"全量用户列表"接口(用户靠群成员反推);C2C 单聊会话通过 + * recentcontact/get_list 按"已同步用户"反查({@link #syncC2CMessages}),漫游消息用 get_roam_msg 全量翻页。 + * 同步进来的消息(群/单聊)均入分发队列(dist_status=0)推送业务系统,本地同时留底。 * * @author imutil */ @@ -33,9 +34,27 @@ public interface SyncService { MigrateTask syncGroupMembers(Long sourceAppId, String tenantId); /** - * 同步群消息:遍历该租户已同步的群,逐群 getGroupMsg + IsFinished 滚动全量 → 写 im_message(source=SYNC) + * 同步群消息:遍历该租户已同步的群,逐群 getGroupMsg + IsFinished 滚动全量 → 写 im_message(source=SYNC,进分发) * * @see #syncGroups(Long, String) 先同步群组,群消息才有目标群清单 */ MigrateTask syncGroupMessages(Long sourceAppId, String tenantId); + + /** + * 同步单聊(C2C)消息:遍历该租户已同步用户 → recentcontact/get_list 收集 C2C 会话 → + * 会话对去重 → get_roam_msg 时间窗向前滚动全量 → 写 im_message(conv_type=1, source=SYNC, 进分发)。 + *
+ * 腾讯无全量 C2C 会话列表接口,靠"已同步用户"反查其最近会话;未导入用户返回 50001 跳过。 + * + * @see #syncGroupMembers(Long, String) 先同步用户(群成员),单聊才有用户清单可遍历 + */ + MigrateTask syncC2CMessages(Long sourceAppId, String tenantId); + + /** + * 补拉对账:重跑群消息 + 单聊同步(幂等,本地已有的跳过、缺失的补入),确保本系统数据最全。 + * 定时任务 {@link com.imutil.task.SyncCheckTask} 驱动,也可手动触发。 + * + * @param tenantId 归属租户(补拉对账仅对主应用归属租户有意义) + */ + MigrateTask checkAndPull(String tenantId); } diff --git a/src/main/java/com/imutil/service/impl/DispatchServiceImpl.java b/src/main/java/com/imutil/service/impl/DispatchServiceImpl.java index 08900ce..4d36178 100644 --- a/src/main/java/com/imutil/service/impl/DispatchServiceImpl.java +++ b/src/main/java/com/imutil/service/impl/DispatchServiceImpl.java @@ -2,9 +2,13 @@ package com.imutil.service.impl; import com.imutil.common.Httpx; import com.imutil.entity.DistQueue; +import com.imutil.entity.ImMessage; +import com.imutil.entity.Tenant; import com.imutil.mapper.DistQueueMapper; import com.imutil.service.DispatchService; +import com.imutil.service.TenantService; import lombok.extern.slf4j.Slf4j; +import org.noear.snack4.ONode; import org.noear.solon.annotation.Component; import org.noear.solon.annotation.Inject; import org.noear.solon.data.annotation.Tran; @@ -29,6 +33,9 @@ public class DispatchServiceImpl implements DispatchService { @Inject private DistQueueMapper distQueueMapper; + @Inject + private TenantService tenantService; + @Inject("${imutil.dispatch.fetchBatch:50}") private int fetchBatch; @@ -91,4 +98,51 @@ public class DispatchServiceImpl implements DispatchService { OffsetDateTime threshold = now.minusMinutes(lockTimeoutMin); return distQueueMapper.recoverStuck(threshold, now); } + + @Override + public void enqueue(String tenantId, ImMessage msg) { + enqueue(tenantId, msg, OffsetDateTime.now()); + } + + /** + * 写入分发队列(复用回调分发链路,payload 为结构化 JSON) + *
+ * 取租户 callbackUrl 作为目标;未配置则跳过(不报错,本地留底仍有效)。
+ * nextRetryAt 由调用方控制:同步批量入队时按序号错开以限速平滑分发。
+ */
+ @Override
+ public void enqueue(String tenantId, ImMessage msg, OffsetDateTime nextRetryAt) {
+ Tenant t = tenantService.getById(tenantId);
+ if (t == null || t.getCallbackUrl() == null || t.getCallbackUrl().isEmpty()) {
+ return;
+ }
+ ONode payload = ONode.ofJson("{}");
+ payload.set("source", msg.getSource());
+ payload.set("msgKey", msg.getMsgKey());
+ payload.set("tenantId", tenantId);
+ payload.set("convType", msg.getConvType());
+ payload.set("convId", msg.getConvId());
+ payload.set("fromAccount", msg.getFromAccount());
+ payload.set("toAccount", msg.getToAccount());
+ payload.set("groupId", msg.getGroupId());
+ payload.set("msgTime", msg.getMsgTime());
+ payload.set("msgType", msg.getMsgType());
+ if (msg.getMsgBody() != null) {
+ try {
+ payload.set("msgBody", ONode.ofJson(msg.getMsgBody()));
+ } catch (Exception e) {
+ payload.set("msgBody", msg.getMsgBody());
+ }
+ }
+ DistQueue q = new DistQueue();
+ q.setMsgKey(msg.getMsgKey());
+ q.setTenantId(tenantId);
+ q.setConvId(msg.getConvId());
+ q.setTargetUrl(t.getCallbackUrl());
+ q.setPayload(payload.toJson());
+ q.setStatus(0);
+ q.setRetryCount(0);
+ q.setNextRetryAt(nextRetryAt);
+ distQueueMapper.insert(q);
+ }
}
diff --git a/src/main/java/com/imutil/service/impl/PullServiceImpl.java b/src/main/java/com/imutil/service/impl/PullServiceImpl.java
index 121a64b..b961991 100644
--- a/src/main/java/com/imutil/service/impl/PullServiceImpl.java
+++ b/src/main/java/com/imutil/service/impl/PullServiceImpl.java
@@ -10,6 +10,7 @@ import com.imutil.entity.Tenant;
import com.imutil.mapper.DistQueueMapper;
import com.imutil.mapper.ImMessageMapper;
import com.imutil.mapper.PullWatermarkMapper;
+import com.imutil.service.DispatchService;
import com.imutil.service.PullService;
import com.imutil.service.TenantService;
import com.imutil.tencent.TencentImClient;
@@ -56,6 +57,9 @@ public class PullServiceImpl implements PullService {
@Inject
private TenantService tenantService;
+ @Inject
+ private DispatchService dispatchService;
+
@Inject("${imutil.pull.convsPerRound:20}")
private int convsPerRound;
@@ -168,50 +172,12 @@ public class PullServiceImpl implements PullService {
imMessageMapper.insert(msg);
newCount++;
// 仅新插入的消息入队分发
- enqueue(wm.getTenantId(), msg);
+ dispatchService.enqueue(wm.getTenantId(), msg);
}
advanceWatermark(wm, maxSeq, maxTime);
return newCount;
}
- /**
- * 写入分发队列(复用回调分发链路,payload 用结构化 JSON)
- */
- private void enqueue(String tenantId, ImMessage msg) {
- Tenant t = tenantService.getById(tenantId);
- if (t == null || t.getCallbackUrl() == null || t.getCallbackUrl().isEmpty()) {
- return;
- }
- ONode payload = ONode.ofJson("{}");
- payload.set("source", msg.getSource());
- payload.set("msgKey", msg.getMsgKey());
- payload.set("tenantId", tenantId);
- payload.set("convType", msg.getConvType());
- payload.set("convId", msg.getConvId());
- payload.set("fromAccount", msg.getFromAccount());
- payload.set("toAccount", msg.getToAccount());
- payload.set("groupId", msg.getGroupId());
- payload.set("msgTime", msg.getMsgTime());
- payload.set("msgType", msg.getMsgType());
- if (msg.getMsgBody() != null) {
- try {
- payload.set("msgBody", ONode.ofJson(msg.getMsgBody()));
- } catch (Exception e) {
- payload.set("msgBody", msg.getMsgBody());
- }
- }
- DistQueue q = new DistQueue();
- q.setMsgKey(msg.getMsgKey());
- q.setTenantId(tenantId);
- q.setConvId(msg.getConvId());
- q.setTargetUrl(t.getCallbackUrl());
- q.setPayload(payload.toJson());
- q.setStatus(0);
- q.setRetryCount(0);
- q.setNextRetryAt(OffsetDateTime.now());
- distQueueMapper.insert(q);
- }
-
/**
* 解析群消息节点为 ImMessage
*/
diff --git a/src/main/java/com/imutil/service/impl/SyncServiceImpl.java b/src/main/java/com/imutil/service/impl/SyncServiceImpl.java
index 1797a9a..225182d 100644
--- a/src/main/java/com/imutil/service/impl/SyncServiceImpl.java
+++ b/src/main/java/com/imutil/service/impl/SyncServiceImpl.java
@@ -15,6 +15,7 @@ import com.imutil.mapper.MigrateTaskMapper;
import com.imutil.mapper.SourceAppMapper;
import com.imutil.mapper.TenantMapper;
import com.imutil.mapper.UserMappingMapper;
+import com.imutil.service.DispatchService;
import com.imutil.service.SyncService;
import com.imutil.tencent.TencentImClient;
import lombok.extern.slf4j.Slf4j;
@@ -25,7 +26,10 @@ import org.noear.solon.annotation.Inject;
import java.time.Instant;
import java.time.OffsetDateTime;
import java.time.ZoneId;
+import java.util.ArrayList;
+import java.util.HashSet;
import java.util.List;
+import java.util.Set;
/**
* 数据同步服务实现(T17)
@@ -89,6 +93,21 @@ public class SyncServiceImpl implements SyncService {
@Inject("${imutil.sync.maxGroups:500}")
private int maxGroups;
+ /** 同步入队分发限速(每秒入队条数,保护业务系统 callbackUrl;错开 next_retry_at 平滑分发) */
+ @Inject("${imutil.sync.dispatchQps:10}")
+ private int syncDispatchQps;
+
+ @Inject
+ private DispatchService dispatchService;
+
+ /** 单聊同步用户上限(防海量用户 N×M 爆配额;如需全量调大) */
+ @Inject("${imutil.sync.maxSyncUsers:200}")
+ private int maxSyncUsers;
+
+ /** 单聊漫游回溯天数(腾讯 C2C 漫游存储期,默认 7 天,套餐更长可调大) */
+ @Inject("${imutil.sync.c2cLookbackDays:7}")
+ private int c2cLookbackDays;
+
/** 数据源密钥解析结果 */
private static class SyncCtx {
long sdkAppId;
@@ -366,6 +385,7 @@ public class SyncServiceImpl implements SyncService {
migrateTaskMapper.updateById(task);
int imported = 0;
+ OffsetDateTime syncBase = OffsetDateTime.now();
for (int gi = 0; gi < groups.size(); gi++) {
String groupId = groups.get(gi).getImGroupId();
try {
@@ -399,6 +419,8 @@ public class SyncServiceImpl implements SyncService {
}
imMessageMapper.insert(msg);
imported++;
+ // 限速入队分发:按序号错开 next_retry_at,DispatchWorker 按速率消费
+ dispatchService.enqueue(ctx.tenantId, msg, nextRetryAt(imported, syncBase));
}
boolean finished = root.get("IsFinished").getLong() == 1;
// 终止:已拉完 / 空页 / 本页无有效 seq
@@ -432,7 +454,7 @@ public class SyncServiceImpl implements SyncService {
}
/**
- * 解析腾讯群消息节点为 ImMessage(source=SYNC,dist_status=1 不进分发)
+ * 解析腾讯群消息节点为 ImMessage(source=SYNC,dist_status=0 进分发)
*/
private ImMessage parseGroupMsgForSync(ONode m, String tenantId, String groupId) {
// 过滤占位/系统消息(IsPlaceMsg=1:From_Account 与 MsgBody 均空,非真实消息,不入库)
@@ -456,10 +478,229 @@ public class SyncServiceImpl implements SyncService {
msg.setMsgBody(body.toString());
msg.setSource("SYNC");
msg.setIsCrossTenant(false);
- msg.setDistStatus(1);
+ msg.setDistStatus(0); // 进分发
return msg;
}
+ // ==================== 同步单聊(C2C) ====================
+
+ /**
+ * 同步单聊消息:遍历已同步用户 → recentcontact/get_list 收集 C2C 会话 → 会话对去重 →
+ * get_roam_msg 时间窗向前滚动全量 → 落库(conv_type=1/source=SYNC/进分发)。
+ */
+ @Override
+ public MigrateTask syncC2CMessages(Long sourceAppId, String tenantId) {
+ SyncCtx ctx = resolveCtx(sourceAppId, tenantId);
+ MigrateTask task = startTask(ctx, "C2C_SYNC");
+ try {
+ List
+ * 定期重跑群消息 + 单聊同步(幂等,本地已有的跳过、缺失的补入),
+ * 兜底腾讯侧有而本系统暂无的消息,确保本系统数据最全、业务系统不漏收。
+ * 对应 app.yml 的 solon.scheduling.job.syncCheckJob。
+ *
+ * 遍历所有租户,对每个用主应用密钥执行 checkAndPull(主应用数据为主)。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class SyncCheckTask {
+
+ @Inject
+ private SyncService syncService;
+
+ @Inject
+ private TenantMapper tenantMapper;
+
+ /**
+ * 由 app.yml syncCheckJob 驱动(默认 fixedDelay=1 小时)
+ */
+ @Scheduled(name = "syncCheckJob")
+ public void run() {
+ List
- * 命令字 openim_admin/get_roam_msg,按时间窗返回指定两人之间最近的 MaxCnt 条消息。
+ * 命令字 openim/admin_getroammsg,按时间窗返回指定两人之间最近的 MaxCnt 条消息。
* 补拉服务据此增量补全回调丢失的单聊消息。
*
* @param fromAccount 发送方 IM 账号
@@ -268,12 +268,13 @@ public class TencentImClient {
*/
public String getRoamMsg(String fromAccount, String toAccount, int maxCnt, long minTime, long maxInterval) {
Map
+ * 命令字 recentcontact/get_list,分页返回 From_Account 的会话 SessionItem:
+ * Type=1 为 C2C(带 To_Account),Type=2 为群(带 GroupId)。
+ * 单聊同步据此发现用户的 C2C 会话对,再配 getRoamMsgRangeAs 全量拉取漫游消息。
+ *
+ * @param fromAccount 本端账号
+ * @param timeStamp 分页时间戳(首轮 0;后续取上轮返回的 TimeStamp)
+ * @param startIndex 分页起始索引(首轮 0;后续取上轮返回的 StartIndex)
+ * @return 腾讯响应原始 JSON(含 SessionItem[]/CompleteFlag/TimeStamp/StartIndex)
+ */
+ public String getRecentContactListAs(String fromAccount, long timeStamp, long startIndex,
+ long srcSdkAppId, String srcSecretKey) {
+ Map
+ * 腾讯 openim/admin_getroammsg 请求字段为 Operator_Account/Peer_Account + MinTime/MaxTime,
+ * 续拉:MaxTime 取上轮返回 LastMsgTime,并带 LastMsgKey,直到 Complete=1。
+ *
+ * @param operatorAccount 会话一方(本端)
+ * @param peerAccount 会话另一方
+ * @param maxCnt 单次条数上限
+ * @param minTime 区间下界(秒级 epoch)
+ * @param maxTime 区间上界(秒级 epoch;首轮 now,续拉取上轮 LastMsgTime)
+ * @param lastMsgKey 上轮返回 LastMsgKey(首轮 null)
+ * @return 腾讯响应原始 JSON(含 MsgList/Complete/LastMsgTime/LastMsgKey)
+ */
+ public String getRoamMsgRangeAs(String operatorAccount, String peerAccount, int maxCnt,
+ long minTime, long maxTime, String lastMsgKey,
+ long srcSdkAppId, String srcSecretKey) {
+ Map 遍历该租户已同步的群,逐群 getGroupMsg + IsFinished 滚动全量 → 写 im_message(source=SYNC,不进分发队列)。请先执行「同步群组」。 遍历该租户已同步的群,逐群 getGroupMsg + IsFinished 滚动全量 → 写 im_message(source=SYNC,进分发队列,限速分批)。请先执行「同步群组」。 遍历该租户已同步用户 → recentcontact/get_list 收集 C2C 会话 → 会话对去重 → get_roam_msg 全量翻页 → 写 im_message(conv_type=1,source=SYNC,进分发)。请先执行「同步群成员」。 重跑群消息 + 单聊同步(幂等:本地已有的跳过、缺失的补入),确保本系统数据最全。仅主应用 + 指定租户,可定时自动或手动触发。同步/迁移任务记录
<#if tasks?has_content>