Browse Source

fix 同步openid、unionid

zoujiajian 3 năm trước cách đây
mục cha
commit
885b84ec38

+ 4 - 0
netflix-dao/src/main/java/com/cyksj/mapper/UserTransferLogMapper.java

@@ -2,6 +2,9 @@ package com.cyksj.mapper;
 
 import com.baomidou.mybatisplus.core.mapper.BaseMapper;
 import com.cyksj.model.entity.UserTransferLog;
+import org.apache.ibatis.annotations.Param;
+
+import java.util.List;
 
 /*
  *项目名: netflix
@@ -10,4 +13,5 @@ import com.cyksj.model.entity.UserTransferLog;
  *创建时间:2023/1/12 10:11
  */
 public interface UserTransferLogMapper extends BaseMapper<UserTransferLog> {
+	void batchInsert(@Param("list") List<UserTransferLog> userTransferLogs);
 }

+ 13 - 0
netflix-dao/src/main/java/com/cyksj/mapper/WxSyncLogMapper.java

@@ -0,0 +1,13 @@
+package com.cyksj.mapper;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.cyksj.model.entity.WxSyncLog;
+
+/*
+ *项目名: netflix
+ *文件名: WxSyncLogMapper
+ *创建者: JavaZou
+ *创建时间:2023/1/12 16:56
+ */
+public interface WxSyncLogMapper extends BaseMapper<WxSyncLog> {
+}

+ 30 - 0
netflix-dao/src/main/java/com/cyksj/model/dto/UserWxDto.java

@@ -0,0 +1,30 @@
+package com.cyksj.model.dto;
+
+import lombok.Getter;
+import lombok.Setter;
+
+import java.util.List;
+
+/*
+ *项目名: netflix
+ *文件名: UserWxDto
+ *创建者: JavaZou
+ *创建时间:2023/1/12 13:46
+ */
+@Getter
+@Setter
+public class UserWxDto {
+	private Integer total;
+
+	private Integer count;
+
+	private Data data;
+
+	private String next_openid;
+
+	@Getter
+	@Setter
+	public static class Data{
+		private List<String> openid;
+	}
+}

+ 20 - 0
netflix-dao/src/main/java/com/cyksj/model/dto/WxOpenIdSyncResultDto.java

@@ -0,0 +1,20 @@
+package com.cyksj.model.dto;
+
+import lombok.Getter;
+import lombok.Setter;
+
+/*
+ *项目名: netflix
+ *文件名: WxOpenIdSyncResultDto
+ *创建者: JavaZou
+ *创建时间:2023/1/12 17:38
+ */
+@Getter
+@Setter
+public class WxOpenIdSyncResultDto {
+	private String ori_openid;
+
+	private String new_openid;
+
+	private String err_msg;
+}

+ 3 - 1
netflix-dao/src/main/java/com/cyksj/model/entity/UserTransferLog.java

@@ -12,7 +12,9 @@ import lombok.Setter;
 @Getter
 @Setter
 public class UserTransferLog extends BaseEntity{
-	private String appId;
+	private String oldAppId;
+
+	private String newAppId;
 
 	private Long userId;
 

+ 30 - 0
netflix-dao/src/main/java/com/cyksj/model/entity/WxSyncLog.java

@@ -0,0 +1,30 @@
+package com.cyksj.model.entity;
+
+import lombok.Getter;
+import lombok.Setter;
+
+/*
+ *项目名: netflix
+ *文件名: WxSyncLog
+ *创建者: JavaZou
+ *创建时间:2023/1/12 16:53
+ */
+@Getter
+@Setter
+public class WxSyncLog extends BaseEntity{
+	private String fromAppid;
+
+	private String toAppid;
+
+
+	private Long userId;
+
+	/**
+	 * 旧公众号openid
+	 */
+	private String oriOpenid;
+
+	private String newOpenid;
+
+	private String errMsg;
+}

+ 13 - 0
netflix-dao/src/main/resources/mapper/UserTransferLogMapper.xml

@@ -0,0 +1,13 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
+<mapper namespace="com.cyksj.mapper.UserTransferLogMapper">
+
+
+    <insert id="batchInsert">
+        insert into user_transfer_log(`old_appid`,`new_appid`,`user_id`,`old_unionid`,`new_unionid`,`type`,`created_time`,`update_time`)
+        values
+        <foreach collection="list" item="item" separator=",">
+            (#{item.oldAppId},#{item.newAppId},#{item.userId},#{item.oldUnionid},#{item.newUnionid},#{item.type},now(),now())
+        </foreach>
+    </insert>
+</mapper>

+ 6 - 0
netflix-service/src/main/java/com/cyksj/config/BaseWeChatConfig.java

@@ -12,6 +12,12 @@ import lombok.Setter;
 @Getter
 @Setter
 public abstract class BaseWeChatConfig {
+
+    /**
+     * 迁移之前的原公众号appid
+     */
+    private String oldAppid;
+
     /**
      * 公众号唯一标识
      */

+ 31 - 13
netflix-service/src/main/java/com/cyksj/service/user/impl/UserServiceImpl.java

@@ -10,14 +10,12 @@ import com.cyksj.common.exception.BusinessRuntimeException;
 import com.cyksj.config.WeChatConfig;
 import com.cyksj.mapper.UserMapper;
 import com.cyksj.mapper.UserTransferLogMapper;
+import com.cyksj.mapper.WxSyncLogMapper;
 import com.cyksj.mapper.manage.distribute.UserDistributeSharedMapper;
 import com.cyksj.mapper.market.task.TaskTypeMapper;
 import com.cyksj.model.dto.GzhOAuth2UserInfo;
 import com.cyksj.model.dto.UserWhoami;
-import com.cyksj.model.entity.TaskType;
-import com.cyksj.model.entity.User;
-import com.cyksj.model.entity.UserDistributeShared;
-import com.cyksj.model.entity.UserTransferLog;
+import com.cyksj.model.entity.*;
 import com.cyksj.service.market.task.TaskService;
 import com.cyksj.service.user.UserService;
 import com.cyksj.service.wechat.WeChatService;
@@ -51,19 +49,34 @@ public class UserServiceImpl extends ServiceImpl<UserMapper, User> implements Us
 
     private final UserTransferLogMapper userTransferLogMapper;
 
+    private final WxSyncLogMapper wxSyncLogMapper;
+
     @Override
     public User saveAuth(GzhOAuth2UserInfo info, Long sharedId) {
         final String unionid = info.getUnionid();
         if (StrUtil.isEmpty(unionid)) {
             throw BusinessRuntimeException.getInstance("授权操作异常,请刷新页面");
         }
-        // 查询用户是否存在
-        User user = Optional.ofNullable(
-                this.getOne(new LambdaQueryWrapper<User>().eq(User::getUnionid,info.getUnionid())))
-                .orElseGet(() -> Optional.ofNullable(this.getOne(new LambdaQueryWrapper<User>().eq(User::getOpenId,info.getOpenid())))
-                        .orElseGet(()-> new User()
-                                .setUnionid(unionid)));
-
+        User user;
+        //迁移公众号之后同步授权后信息
+        //是否是迁移用户
+        WxSyncLog wxSyncLog = wxSyncLogMapper.selectOne(Wrappers.lambdaQuery(WxSyncLog.class)
+                .eq(WxSyncLog::getToAppid, weChatConfig.getAppid())
+                .eq(WxSyncLog::getNewOpenid, info.getOpenid()).last("limit 1"));
+        if (wxSyncLog != null) {
+            user = this.getOne(Wrappers.lambdaQuery(User.class)
+                    .eq(User::getOpenId, wxSyncLog.getOriOpenid()).last("limit 1"));
+            if (user != null) {
+                user.setOpenId(wxSyncLog.getNewOpenid());
+            }
+        } else {
+            // 查询用户是否存在
+            user = Optional.ofNullable(
+                    this.getOne(new LambdaQueryWrapper<User>().eq(User::getUnionid, info.getUnionid())))
+                    .orElseGet(() -> this.getOne(new LambdaQueryWrapper<User>().eq(User::getOpenId, info.getOpenid())));
+        }
+        user = Optional.ofNullable(user).orElseGet(() -> new User()
+                .setUnionid(unionid));
         user.setNickname(info.getNickname())
                 .setSex(info.getSex())
                 .setCountry(info.getCountry())
@@ -72,9 +85,14 @@ public class UserServiceImpl extends ServiceImpl<UserMapper, User> implements Us
                 .setHeadimgurl(info.getHeadimgurl());
         if (!user.getUnionid().equals(unionid)) {
             UserTransferLog log = new UserTransferLog();
-            log.setAppId(weChatConfig.getAppid());
+            if (wxSyncLog != null) {
+                log.setType("公众号迁移");
+            } else {
+                log.setType("绑定开发者账号");
+            }
+            log.setOldAppId(weChatConfig.getOldAppid());
+            log.setNewAppId(weChatConfig.getAppid());
             log.setUserId(user.getId());
-            log.setType("授权同步");
             log.setOldUnionid(user.getUnionid());
             log.setNewUnionid(info.getUnionid());
             userTransferLogMapper.insert(log);

+ 215 - 45
netflix-web/src/main/java/com/cyksj/web/controller/user/UserController.java

@@ -1,16 +1,23 @@
 package com.cyksj.web.controller.user;
 
+import cn.hutool.core.collection.CollUtil;
 import cn.hutool.core.date.DateTime;
 import cn.hutool.core.util.StrUtil;
+import cn.hutool.json.JSONObject;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
 import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import com.cyksj.common.util.J11HttpC;
+import com.cyksj.common.util.Jsons;
+import com.cyksj.config.WeChatConfig;
 import com.cyksj.dto.Result;
 import com.cyksj.enums.GatewayResponse;
 import com.cyksj.mapper.OrderDonMapper;
 import com.cyksj.mapper.UserTransferLogMapper;
+import com.cyksj.mapper.WxSyncLogMapper;
 import com.cyksj.mapper.manage.coupon.CouponUserMapper;
 import com.cyksj.mapper.market.task.UserBindPhoneMapper;
-import com.cyksj.model.dto.GzhUnionidUserinfo;
-import com.cyksj.model.dto.UserWhoami;
+import com.cyksj.model.dto.*;
 import com.cyksj.model.entity.*;
 import com.cyksj.model.request.ExchangeCodeTicket;
 import com.cyksj.model.views.RenewalView;
@@ -20,11 +27,18 @@ import com.cyksj.service.relation.GroupRelationFrontService;
 import com.cyksj.service.user.UserService;
 import com.cyksj.service.wechat.WeChatService;
 import com.cyksj.web.util.StpUserUtil;
+import com.google.common.collect.Lists;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.web.bind.annotation.*;
 
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.util.ArrayList;
+import java.util.Date;
 import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
 
 /**
  * @author chan
@@ -106,47 +120,203 @@ public class UserController {
 
     private final UserTransferLogMapper userTransferLogMapper;
 
-    /**
-     * 获取微信用户信息
-     */
-    @GetMapping("/get/wx/user")
-    public void syncUser() {
-        List<User> users = userService.list(Wrappers.lambdaQuery(User.class)
-                .ne(User::getUnionid, ""));
-        String appId = "wxd58016eee53b6558";
-        users.forEach(user->{
-            try {
-                String accessToken = "64_rPJYk5-812w8P9pMMn-X4JYkuxg-NFHp6TNWwdFl1KgTkEhwPXGPrSl624mmD5WebnG0Ll37-iQ8bOnUQtOE82DLj1pkRC1zvEEnXhEQh_wbEwak-ijroBeSqpcZMOaAJAMZB";
-                GzhUnionidUserinfo userInfo2 = weChatService.getUserInfo2(accessToken, user.getOpenId());
-                if ("ADD_SCENE_ACCOUNT_MIGRATION".equals(userInfo2.getSubscribeScene())) {
-                    log.info(user.getOpenId() + "公众号迁移");
-                }
-                if (userInfo2 == null) {
-                    log.info(user.getOpenId() + "失联");
-                }
-                if (StrUtil.isEmpty(userInfo2.getUnionid())) {
-                    log.info(user.getOpenId() + "未关注新公众号");
-                    return;
-                }
-                if (user.getUnionid().equals(userInfo2.getUnionid())) {
-                    log.info("新旧公众号unionid相同");
-                    return;
-                } else {
-                    user.setUnionid(userInfo2.getUnionid());
-                    userService.updateById(user);
-
-                    UserTransferLog log = new UserTransferLog();
-                    log.setAppId(appId);
-                    log.setUserId(user.getId());
-                    log.setOldUnionid(user.getUnionid());
-                    log.setType("手动同步");
-                    log.setNewUnionid(userInfo2.getUnionid());
-
-                    userTransferLogMapper.insert(log);
-                }
-            } catch (Exception e) {
-                System.out.println("获取accesstoken错误");
-            }
-        });
-    }
+	private final WxSyncLogMapper wxSyncLogMapper;
+
+	private final WeChatConfig weChatConfig;
+
+	/**
+	 * 获取微信用户信息, 迁移后通过openid同步用户unionid
+	 */
+	@GetMapping("/get/wx/user")
+	public void syncUser() {
+		String oldAppId = weChatConfig.getOldAppid();
+		String newAppId = weChatConfig.getAppid();
+		LambdaQueryWrapper<WxSyncLog> wrapper = Wrappers.lambdaQuery(WxSyncLog.class)
+				.ne(WxSyncLog::getNewOpenid, "")
+				.in(WxSyncLog ::getUserId,List.of(9909,11201))
+				.orderByAsc(WxSyncLog::getId);
+		Integer limit = 12000;
+		Integer count = wxSyncLogMapper.selectCount(wrapper);
+		int num = (count / limit) + 1;
+		for (int i = 0; i < num; i++) {
+			Integer start = i + 1;
+			Page<WxSyncLog> wxSyncLogs = wxSyncLogMapper.selectPage(new Page<>(start, limit),wrapper);
+			List<WxSyncLog> records = wxSyncLogs.getRecords();
+			List<UserTransferLog> userTransferLogs = new ArrayList<>(records.size());
+			records.forEach(syncUser -> {
+				try {
+					String accessToken = weChatService.getAccessToken();
+					String oriOpenid = syncUser.getOriOpenid();
+					String newOpenId = syncUser.getNewOpenid();
+					GzhUnionidUserinfo userInfo2 = weChatService.getUserInfo2(accessToken, newOpenId);
+					if (userInfo2 == null) {
+						log.info("newopenid:{}未关注新公众号", newOpenId);
+						return;
+					}
+					if ("ADD_SCENE_ACCOUNT_MIGRATION".equals(userInfo2.getSubscribeScene())) {
+						log.info("旧oriOpenid:{}->新newOpenid:{}公众号迁移", oriOpenid, newOpenId);
+					}
+					if (userInfo2 == null) {
+						log.info("该newopenid:{}不是该新公众号的用户", newOpenId);
+						return;
+					}
+					User user = userService.getOne(Wrappers.lambdaQuery(User.class)
+							.eq(User::getOpenId, oriOpenid)
+							.select(User::getId,
+									User::getUnionid,
+									User::getOpenId).last("limit 1"));
+					if (user == null) {
+						log.info("该newopenid:{}是新公众号的新用户,同步过来的数据不存在这种情况", newOpenId);
+						return;
+					}
+					if (StrUtil.isEmpty(userInfo2.getUnionid())) {
+						log.info("该用户未未关注新公众号,等待在新公众号授权重置其uinonid");
+					}
+					if (user.getUnionid() != null && user.getUnionid().equals(userInfo2.getUnionid())) {
+						log.info("用户id:{},新旧公众号unionid相同", user.getId());
+						return;
+					} else {
+						String unionid = user.getUnionid();
+						user.setOpenId(newOpenId);
+						if (StrUtil.isNotEmpty(userInfo2.getUnionid())) {
+							user.setUnionid(userInfo2.getUnionid());
+						}
+//						userService.updateById(user);
+
+						UserTransferLog log = new UserTransferLog();
+						log.setOldAppId(oldAppId);
+						log.setNewAppId(newAppId);
+						log.setUserId(user.getId());
+						log.setOldUnionid(unionid);
+						log.setType("手动同步");
+						log.setNewUnionid(userInfo2.getUnionid());
+						userTransferLogs.add(log);
+					}
+				} catch (Exception e) {
+					System.out.println("获取accesstoken错误");
+				}
+			});
+			List<List<UserTransferLog>> logsList = Lists.partition(userTransferLogs, 100);
+			logsList.forEach(insertLogList->{
+				userTransferLogMapper.batchInsert(insertLogList);
+			});
+			userTransferLogs.clear();
+		}
+	}
+
+	/**
+	 * 同步旧的公众号openid->新公众号
+	 */
+	@GetMapping("/sync/old/wx/user")
+	public void test() throws Exception {
+		String accessToken = "+ \"64_F1RV0DPgP6silfzERCbrckClRYLk2-HsP2mklVETA_t2FRlnlnN0wpBcFcy5rBqMv5VPpwwNsSM47Y1uuuoEDkgtWQLtBEUH3EgtnNmMAZYa8brtKvf_GRTJiPcHKFbAIAQTO\"";
+		String syncUrl = "https://api.weixin.qq.com/cgi-bin/changeopenid?access_token=" + weChatService.getAccessToken();
+//		UserWxDto userWxDto = getUserWxDto("");
+//		syncWxUser(userWxDto, syncUrl);
+		Integer limit = 12000;
+		WxSyncLog wxSyncLog = wxSyncLogMapper.selectOne(Wrappers.lambdaQuery(WxSyncLog.class)
+				.orderByDesc(WxSyncLog::getId)
+				.last("limit 1"));
+		String oriOpenId = wxSyncLog.getOriOpenid();
+		Date finalCreateTime = userService.getOne(Wrappers.lambdaQuery(User.class)
+				.eq(User::getOpenId, oriOpenId).last("limit 1")).getCreatedTime();
+		LambdaQueryWrapper<User> select = Wrappers.lambdaQuery(User.class)
+				.ne(User::getOpenId, "")
+				.ge(User::getCreatedTime, finalCreateTime)
+				.ne(User::getOpenId, oriOpenId)
+				.orderByAsc(User::getId)
+				.select(User::getOpenId);
+		int count = userService.count(select);
+		int num = (count / limit) + 1;
+		for (int i = 0; i < num; i++) {
+			Integer start = i + 1;
+			Page<User> page = userService.page(new Page<>(start, limit), select.select(User::getOpenId, User::getId));
+			List<User> records = page.getRecords();
+			syncWxUser(records, syncUrl);
+		}
+	}
+
+	private void syncWxUser(List<User> users, String syncUrl) throws Exception {
+		List<String> openidList = users.stream().map(User::getOpenId).collect(Collectors.toList());
+		Map<String, List<User>> openidMap = users.stream().collect(Collectors.groupingBy(User::getOpenId));
+		if (CollUtil.isNotEmpty(openidList)) {
+			List<List<String>> partition = Lists.partition(openidList, 100);
+			partition.forEach(openids -> {
+				JSONObject jsonObject = new JSONObject();
+				jsonObject.putOpt("from_appid", weChatConfig.getOldAppid());
+				jsonObject.putOpt("openid_list", openids);
+				try {
+					HttpResponse<String> send = J11HttpC.custom()
+							.ofPost()
+							.url(syncUrl)
+							.body(HttpRequest.BodyPublishers.ofString(jsonObject.toString()))
+							.send(HttpResponse.BodyHandlers.ofString());
+
+					JSONObject json = Jsons.parseObject(send.body(), JSONObject.class);
+					List<WxOpenIdSyncResultDto> result_list = Jsons.parseList(json.get("result_list"), WxOpenIdSyncResultDto.class);
+					result_list.forEach(result -> {
+						WxSyncLog wxSyncLog = new WxSyncLog();
+						wxSyncLog.setFromAppid(weChatConfig.getOldAppid());
+						wxSyncLog.setToAppid(weChatConfig.getAppid());
+						wxSyncLog.setUserId(openidMap.get(result.getOri_openid()).get(0).getId());
+						wxSyncLog.setOriOpenid(result.getOri_openid());
+						wxSyncLog.setNewOpenid(result.getNew_openid());
+						wxSyncLog.setErrMsg(result.getErr_msg());
+						wxSyncLogMapper.insert(wxSyncLog);
+					});
+				} catch (Exception e) {
+					log.error("公众号迁移错误,error_msg:", e.getMessage());
+				}
+//				List<User> user = userService.list(Wrappers.lambdaQuery(User.class)
+//						.in(User::getOpenId, openids));
+//				if (user.isEmpty()) {
+//					log.info("ss");
+//				}
+//				if (user.size() != openids.size()) {
+//					log.info("不相等");
+//					List<String> collect = user.stream().map(User::getOpenId).collect(Collectors.toList());
+//					log.info("本地用户大小:{},微信关注用户数:{}", user.size(), openids.size());
+//					openids.removeAll(collect);
+//					try {
+//						log.info("剩余:{}", openids);
+//					} catch (Exception e) {
+//						e.printStackTrace();
+//					}
+//				}
+//				});
+//			});
+//			if (StrUtil.isNotBlank(userWxDto.getNext_openid())) {
+//				userWxDto = getUserWxDto(userWxDto.getNext_openid());
+//				syncWxUser(userWxDto, syncUrl);
+//			}
+			});
+		}
+	}
+
+
+	private UserWxDto getUserWxDto(String nxt_openid) throws Exception {
+		String accessToken = "64_F1RV0DPgP6silfzERCbrckClRYLk2-HsP2mklVETA_t2FRlnlnN0wpBcFcy5rBqMv5VPpwwNsSM47Y1uuuoEDkgtWQLtBEUH3EgtnNmMAZYa8brtKvf_GRTJiPcHKFbAIAQTO";
+		String url = "https://api.weixin.qq.com/cgi-bin/user/get?access_token=%s&next_openid=%s";
+		HttpResponse<String> response = J11HttpC.custom()
+				.ofGet()
+				.url(String.format(url, accessToken, nxt_openid))
+				.send(HttpResponse.BodyHandlers.ofString());
+		return Jsons.parseObject(response.body(), UserWxDto.class);
+	}
+
+	@GetMapping("/test")
+	public void tes2t() throws Exception {
+//		List<String> strings = List.of("oXxcY6I5boH-4sSBf3eL5lARrGP0", "oXxcY6FY-PFSRohN-NuiIb_XKHXM","oXxcY6A5GdZIPbZBj4CFkXcNWVjU");
+//		strings.forEach(opneid->{
+//			try {
+//				GzhUnionidUserinfo userInfo2 = weChatService.getUserInfo2(weChatService.getAccessToken(), opneid);
+//				System.out.println(Jsons.toJson(userInfo2));
+//			} catch (Exception e) {
+//				e.printStackTrace();
+//			}
+//		});
+		AuthToken authToken = weChatService.getAuthToken("081UY51w3Q7AWZ2dlc0w3EJPFp4UY51j");
+		GzhOAuth2UserInfo userInfo = weChatService.getUserInfo(authToken.getAccessToken(), authToken.getOpenid());
+		System.out.println("s");
+	}
 }

Những thai đổi đã bị hủy bỏ vì nó quá lớn
+ 0 - 0
netflix-web/src/main/resources/application-dev.yml


Những thai đổi đã bị hủy bỏ vì nó quá lớn
+ 0 - 0
netflix-web/src/main/resources/application-prd.yml


+ 1 - 1
netflix-web/src/main/resources/application.yml

@@ -1 +1 @@
-spring:
  profiles:
    active: prd
    
+spring:
  profiles:
    active: dev
    

Một số tệp đã không được hiển thị bởi vì quá nhiều tập tin thay đổi trong này khác