Просмотр исходного кода

裂变计划redis v0.06 +听书 +其他

MOKASZ\lw12420 2 лет назад
Родитель
Сommit
2021c87c6a
19 измененных файлов с 1159 добавлено и 21 удалено
  1. 2 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/controller/planNew/PlanNewCommonOutController.java
  2. 1 1
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/controller/planNew/PlanNewOutController.java
  3. 16 6
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/controller/testController.java
  4. 5 5
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/service/listening/ListeningPlanServiceImpl.java
  5. 30 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/MatchByKeyDeleteAll.java
  6. 193 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncBackgroupGroupImageItemToRedis.java
  7. 195 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncGroupUsersToRedis.java
  8. 209 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncImageToRedis.java
  9. 1 1
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncListeningAudioToRedis.java
  10. 202 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncMaterialToRedis.java
  11. 3 3
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncPlanNewListenToRedis.java
  12. 4 4
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncPlanNewToRedis.java
  13. 214 0
      launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncReactiveDomainToRedis.java
  14. 1 1
      launch-admin/src/test/java/com/webflux/launchadmin/LaunchAdminApplicationTests.java
  15. 15 0
      launch-redis/src/main/java/com/weblux/launchredis/utils/BackgroupImageItemRedisKey.java
  16. 17 0
      launch-redis/src/main/java/com/weblux/launchredis/utils/GroupUsersRedisKey.java
  17. 17 0
      launch-redis/src/main/java/com/weblux/launchredis/utils/ImageRedisKey.java
  18. 17 0
      launch-redis/src/main/java/com/weblux/launchredis/utils/MaterialRedisKey.java
  19. 17 0
      launch-redis/src/main/java/com/weblux/launchredis/utils/ReactiveDomainRedisKey.java

+ 2 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/controller/planNew/PlanNewCommonOutController.java

@@ -73,6 +73,7 @@ public class PlanNewCommonOutController {
 
 
 
+
     private static final String COLUMN="id,name,code,type,poster_template_id,group_id,config_json,import_domain,preview_url";
     private static final String[] COLUMNARRAYCOMMON=COLUMN.split(",");
     /**
@@ -181,6 +182,7 @@ public class PlanNewCommonOutController {
      */
     @GetMapping("getPlanNewCommonH5")
     public Mono<RStatus<byte[]>> getPlanNewCommonH5(@RequestParam("code") String code, ServerWebExchange exchange) {
+
         Mono<PlanNewCommon> planNewCommonMono = template.selectOne(Query.query(Criteria.where("code").is(code)).columns(COLUMNARRAYCOMMON), PlanNewCommon.class)
                 .switchIfEmpty(Mono.error(new BaseException("计划code异")))
                 .onErrorResume(throwable -> Mono.error(new BaseException("计划code异")));

+ 1 - 1
launch-admin/src/main/java/com/webflux/launchadmin/mysql/controller/planNew/PlanNewOutController.java

@@ -72,7 +72,7 @@ public class PlanNewOutController {
     @Resource
     private PlanServiceInterface planServiceInterface;
 
-    private static final ConcurrentHashMap<String,Integer> map = new ConcurrentHashMap<String,Integer>();
+    //private static final ConcurrentHashMap<String,Integer> map = new ConcurrentHashMap<String,Integer>();
 
     //清空用户Map数据
     // @Scheduled(cron="1 1 4 * * ?")

+ 16 - 6
launch-admin/src/main/java/com/webflux/launchadmin/mysql/controller/testController.java

@@ -8,8 +8,10 @@ import com.webflux.launchadmin.context.UserContext;
 import com.webflux.launchadmin.context.XxlSsoUser;
 import com.webflux.launchadmin.mysql.entity.listening.ListeningAudio;
 import com.webflux.launchadmin.mysql.entity.listening.ListeningPlanNew;
+import com.webflux.launchadmin.mysql.entity.planNew.Image;
 import com.webflux.launchadmin.mysql.entity.test.Tweet;
 import com.webflux.launchadmin.mysql.repository.test.TweetRepository;
+import com.webflux.launchadmin.mysql.task.SyncImageToRedis;
 import com.webflux.launchadmin.mysql.task.SyncListeningAudioToRedis;
 import com.webflux.launchadmin.mysql.task.SyncPlanNewListenToRedis;
 import com.webflux.launchcommon.returnObj.Paged;
@@ -58,17 +60,25 @@ public class testController {
     @Resource
     private R2dbcEntityTemplate template;
     @Resource
-    private SyncListeningAudioToRedis redis;
-
+    private SyncImageToRedis redis;
+    @Resource
+    private SyncListeningAudioToRedis redis1;
     @GetMapping("get5")
     public Mono<List<ListeningAudio>>  get5()   {
-        return redis.byRedisGetLinteningAudioIfErByMysql(3);
+        redis1.syncLinteningAudioPlan();
+        return redis1.byRedisGetLinteningAudioIfErByMysql(3);
     }
-    @GetMapping("get4")
-    public Mono<Void>  get4()   {
-        //return syncPlanNewListenToRedis.syncLinteningPlan();
+    @GetMapping("get8")
+    public Mono<Void>  get8()   {
+        redis1.deleteAll();
         return Mono.empty();
     }
+    @GetMapping("get4")
+    public Mono<List<Image>>  get4()   {
+        redis.syncImagePlan();
+       return redis.byRedisGetImageIfErByMysql("19","10");
+
+    }
 
 
     @GetMapping("get3")

+ 5 - 5
launch-admin/src/main/java/com/webflux/launchadmin/mysql/service/listening/ListeningPlanServiceImpl.java

@@ -240,7 +240,7 @@ public class ListeningPlanServiceImpl implements ListeningPlanServiceInterface{
                                     stringBuilder.append("&sk=");
                                     stringBuilder.append(IdUtil.fastSimpleUUID());
                                 }else {
-                                    throw new BaseException("通过计划id:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
+                                    throw new BaseException("通过计划id3:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
                                 }
                                 Optional<Item> first = collect.stream()
                                         .filter(ft -> ft.getPlanId()
@@ -575,7 +575,7 @@ public class ListeningPlanServiceImpl implements ListeningPlanServiceInterface{
                         stringBuilder.append("&sk=");
                         stringBuilder.append(IdUtil.fastSimpleUUID());
                     }else {
-                        throw new BaseException("二维码链接:::通过计划id:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
+                        throw new BaseException("二维码链接:::通过计划id4:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
                     }
                     return stringBuilder;
                 })
@@ -822,7 +822,7 @@ public class ListeningPlanServiceImpl implements ListeningPlanServiceInterface{
                                     stringBuilder.append("&sk=");
                                     stringBuilder.append(IdUtil.fastSimpleUUID());
                                 }else {
-                                    throw new BaseException("通过计划id:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
+                                    throw new BaseException("通过计划id5:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
                                 }
                                 Optional<Item> first = collect.stream()
                                         .filter(ft -> ft.getPlanId()
@@ -1180,7 +1180,7 @@ public class ListeningPlanServiceImpl implements ListeningPlanServiceInterface{
             stringBuilder.append("&sk=");
             stringBuilder.append(IdUtil.fastSimpleUUID());
         }else {
-            throw new BaseException("通过计划id:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
+            throw new BaseException("通过计划id1:"+a.listeningPlanNewId() +"的营销类型查询域名为空");
         }
     }
     private void getUrl(PlanReq saveRequest, List<DomainSelect> b, StringBuilder stringBuilder) {
@@ -1195,7 +1195,7 @@ public class ListeningPlanServiceImpl implements ListeningPlanServiceInterface{
             stringBuilder.append("&sk=");
             stringBuilder.append(IdUtil.fastSimpleUUID());
         }else {
-            throw new BaseException("通过计划id:"+saveRequest.getPlanNewId() +"的营销类型查询域名为空");
+            throw new BaseException("通过计划id2:"+saveRequest.getPlanNewId() +"的营销类型查询域名为空");
         }
     }
 

+ 30 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/MatchByKeyDeleteAll.java

@@ -0,0 +1,30 @@
+package com.webflux.launchadmin.mysql.task;
+
+import jakarta.annotation.Resource;
+import org.springframework.data.redis.core.ReactiveRedisTemplate;
+import org.springframework.data.redis.core.ScanOptions;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+public class MatchByKeyDeleteAll {
+    @Resource
+    private ReactiveRedisTemplate<String,String> reactiveRedisTemplate;
+    public Mono<Void> deleteAll() {
+        deleteByPattern("sync_image_*").subscribe();
+        deleteByPattern("sync_lintening_*").subscribe();
+        deleteByPattern("sync_plan_new_*").subscribe();
+        //   reactiveRedisTemplate.delete("sync_lintening_*").subscribe();
+        return Mono.empty();
+    }
+    public Mono<Long> deleteByPattern(String pattern) {
+        return reactiveRedisTemplate
+                .scan(ScanOptions.scanOptions().match(pattern).build())
+                .collectList()
+                .flatMap(keys -> {
+                    if (keys.isEmpty()) {
+                        return Mono.just(0L);
+                    }
+                    return reactiveRedisTemplate.delete(Flux.fromIterable(keys));
+                });
+    }
+}

+ 193 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncBackgroupGroupImageItemToRedis.java

@@ -0,0 +1,193 @@
+package com.webflux.launchadmin.mysql.task;
+
+import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationContext;
+import com.fasterxml.jackson.databind.JsonDeserializer;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import com.fasterxml.jackson.databind.module.SimpleModule;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
+import com.plumelog.core.util.LogExceptionStackTrace;
+import com.webflux.launchadmin.global.BaseException;
+import com.webflux.launchadmin.mysql.entity.groupUsers.GroupUsers;
+import com.webflux.launchadmin.mysql.entity.planNew.BackgroupImageItem;
+import com.weblux.launchredis.utils.BackgroupImageItemRedisKey;
+import com.weblux.launchredis.utils.GroupUsersRedisKey;
+import jakarta.annotation.Resource;
+import lombok.extern.slf4j.Slf4j;
+import net.javacrumbs.shedlock.spring.annotation.SchedulerLock;
+import org.jetbrains.annotations.NotNull;
+import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
+import org.springframework.data.redis.core.ReactiveRedisTemplate;
+import org.springframework.data.relational.core.query.Criteria;
+import org.springframework.data.relational.core.query.Query;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+import java.io.IOException;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+@Component
+@Slf4j
+public class SyncBackgroupGroupImageItemToRedis {
+    @Resource
+    private R2dbcEntityTemplate template;
+    @Resource
+    private ReactiveRedisTemplate<String,String> reactiveRedisTemplate;
+
+    private static final String COLUMN="id,image,index,group_id";
+    private static final String[] COLUMNARRAY=COLUMN.split(",");
+    /**
+     *同步裂变计划
+     */
+    @Scheduled(fixedDelay = 100000)
+    @SchedulerLock(name = "launch:SyncBackgroupGroupImageItemToRedis:syncBackgroupItemPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
+    public void syncBackgroupItemPlan(){
+        System.out.println("start");
+        Mono<List<BackgroupImageItem>> listMono = template
+                .select(Query.query(Criteria.empty())
+                        .columns(COLUMNARRAY),BackgroupImageItem.class)
+                .collectList();
+          listMono.flatMap(list -> {
+              var collect1 = list.stream().collect(Collectors.groupingBy(BackgroupImageItem::groupId));
+              collect1.forEach(((x,y)->
+              {
+                  var groupUsersIdMysql = y.stream()
+                          .map(BackgroupImageItem::groupId)
+                          .map(String::valueOf).collect(Collectors.toList());
+                  reactiveRedisTemplate
+                          .opsForSet()
+                          .members(BackgroupImageItemRedisKey.BACKGROUP_IMAGE_ITEM_INFO_SET+x)
+                          .switchIfEmpty(Mono.just(""))
+                          .collectList().flatMap(listIds -> {
+                      var addList = groupUsersIdMysql.stream().filter(f->!listIds.contains(f)).collect(Collectors.toList());
+                      if (!addList.isEmpty()) {
+                          var collect = y.stream()
+                                  .filter(f -> addList.contains(String.valueOf(f.groupId())))
+                                  .collect(Collectors.toList());
+                          collect.forEach(groupUser-> setRedisValue(groupUser, String.valueOf(x)).subscribe());
+                      }
+                      var deleteList = listIds.stream().filter(f->!groupUsersIdMysql.contains(f)).collect(Collectors.toList());
+                      if (!deleteList.isEmpty()) {
+                          deleteList.forEach(id-> executeDelete(id, String.valueOf(x)).subscribe());
+                      }
+                      return Mono.empty();
+                  }).subscribe();
+              }
+              ));
+              return Mono.empty();
+        }).subscribe();
+        System.out.println("end");
+    }
+
+    @NotNull
+    public Mono<Void> setRedisValue(BackgroupImageItem planNew,String key) {
+        reactiveRedisTemplate.opsForSet().add(BackgroupImageItemRedisKey.BACKGROUP_IMAGE_ITEM_INFO_SET+key, String.valueOf(planNew.groupId())).subscribe();
+        String json= getPlanNewToString(planNew);
+        reactiveRedisTemplate.opsForHash().put(BackgroupImageItemRedisKey.BACKGROUP_IMAGE_ITEM_INFO_HASH,String.valueOf(planNew.id()), json).subscribe();
+       return Mono.empty();
+    }
+    @NotNull
+    public Mono<Void> executeDelete(String id,String key) {
+        reactiveRedisTemplate.opsForSet().remove(BackgroupImageItemRedisKey.BACKGROUP_IMAGE_ITEM_INFO_SET+key, id).subscribe();
+        reactiveRedisTemplate.opsForHash().remove(BackgroupImageItemRedisKey.BACKGROUP_IMAGE_ITEM_INFO_HASH , id ).subscribe();
+        return Mono.empty();
+    }
+
+    private String getPlanNewToString(BackgroupImageItem planNew) {
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
+        try {
+            // 将 record 转换为 JSON 字符串
+            return objectMapper.writeValueAsString(planNew);
+        } catch (JsonProcessingException e) {
+           throw new BaseException("PlanNew 将 record 转换为 JSON 字符串 异常:::"+e.getMessage());
+        }
+    }
+
+    /**
+     * 通过code从数据库获取PlanNew并缓存redis
+     * @param backgroundGroupId
+     * @return
+     */
+    public Mono<List<BackgroupImageItem>> byCodeAndMysqlGetBackgroupImageItem(String backgroundGroupId){
+        Flux<BackgroupImageItem> backgroupImageItemFlux = template
+                .select(Query.query(Criteria.where("group_id")
+                        .is(backgroundGroupId)), BackgroupImageItem.class)//Mono.error(new BaseException("查询海报模版 id:"+f.posterTemplateId()+"=>空"))
+                .onErrorResume(throwable -> Mono.error(new BaseException("byCodeAndMysqlGetBackgroupImageItem:::"+throwable.getMessage())));
+        return backgroupImageItemFlux.flatMap(planNew -> {
+          setRedisValue(planNew,backgroundGroupId).subscribe();
+           return Mono.just(planNew);
+        }).collectList();
+    }
+
+    /**
+     * 通过redis获取缓存 空或失败 从mysql中取
+     * @param backgroundGroupId
+     * @return
+     */
+    public Mono<List<BackgroupImageItem>> byRedisGetPlanNewIfErByMysql(String backgroundGroupId){
+      return   reactiveRedisTemplate.opsForSet()
+                 .members(BackgroupImageItemRedisKey.BACKGROUP_IMAGE_ITEM_INFO_SET+backgroundGroupId)
+              .collectList()
+                 .switchIfEmpty(Mono.just(List.of()))
+              .onErrorResume(throwable -> {
+                  log.error("接口异常  {} e = {}","获取redis BackgroupImageItem 失败 暂时从Mysql获取内容" , LogExceptionStackTrace.erroStackTrace(throwable));
+                  return Mono.just(List.of());
+              })
+              .flatMap(list -> {
+                  if( Objects.nonNull(list) && !list.isEmpty()){
+                      return   Flux.fromIterable(list).flatMap(f-> reactiveRedisTemplate.opsForHash()
+                              .get(GroupUsersRedisKey.GROUP_USERS_INFO_HASH , f))
+                              .flatMap(fm->{
+                          var backgroupImageItem = strParsePlanNew(String.valueOf(fm));
+                          return Mono.just(backgroupImageItem);
+                      }).collectList();
+                  }else {
+                      return byCodeAndMysqlGetBackgroupImageItem(backgroundGroupId);
+                  }
+         });
+    }
+    /**
+     * str 转PlanNew
+     * @param str
+     * @return
+     */
+    public BackgroupImageItem strParsePlanNew(String str) {
+        // 创建 ObjectMapper 实例并配置
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        SimpleModule module = new SimpleModule();
+        module.addDeserializer(ZonedDateTime.class, new CustomZonedDateTimeDeserializer());
+        objectMapper.registerModule(module);
+            try {
+                // 将 JSON 字符串转换为 record
+                return objectMapper.readValue(str, BackgroupImageItem.class);
+            } catch (JsonProcessingException e) {
+                throw new BaseException(" 将 JSON字符串 转换为 BackgroupImageItem  异常:::"+e.getMessage());
+            }
+    }
+    // 自定义反序列化器
+    public static class CustomZonedDateTimeDeserializer extends JsonDeserializer<ZonedDateTime> {
+        private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+        @Override
+        public ZonedDateTime deserialize(JsonParser p, DeserializationContext ctxt)
+                throws IOException, JsonProcessingException {
+            String date = p.getText();
+            LocalDateTime localDateTime = LocalDateTime.parse(date, FORMATTER);
+            // 使用系统默认时区,可以根据需要修改
+            return ZonedDateTime.of(localDateTime, ZoneId.systemDefault());
+        }
+    }
+}

+ 195 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncGroupUsersToRedis.java

@@ -0,0 +1,195 @@
+package com.webflux.launchadmin.mysql.task;
+
+import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationContext;
+import com.fasterxml.jackson.databind.JsonDeserializer;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import com.fasterxml.jackson.databind.module.SimpleModule;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
+import com.plumelog.core.util.LogExceptionStackTrace;
+import com.webflux.launchadmin.global.BaseException;
+import com.webflux.launchadmin.mysql.entity.groupUsers.GroupUsers;
+import com.webflux.launchadmin.mysql.entity.planNew.Image;
+import com.weblux.launchredis.utils.GroupUsersRedisKey;
+import jakarta.annotation.Resource;
+import lombok.extern.slf4j.Slf4j;
+import net.javacrumbs.shedlock.spring.annotation.SchedulerLock;
+import org.jetbrains.annotations.NotNull;
+import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
+import org.springframework.data.redis.core.ReactiveRedisTemplate;
+import org.springframework.data.relational.core.query.Criteria;
+import org.springframework.data.relational.core.query.Query;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+import java.io.IOException;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+@Component
+@Slf4j
+public class SyncGroupUsersToRedis {
+    @Resource
+    private R2dbcEntityTemplate template;
+    @Resource
+    private ReactiveRedisTemplate<String,String> reactiveRedisTemplate;
+
+    private static final String COLUMN="group_users_id,group_users_head,group_users_nickname,group_users_type";
+    private static final String[] COLUMNARRAY=COLUMN.split(",");
+
+//    Mono<List<GroupUsers>> groupUsersMono = template.select(Query.query(Criteria
+//            .where("group_users_type").is(f.groupId())).limit(14), GroupUsers.class)
+//            .onErrorResume(throwable -> Mono.error(new BaseException("获取群成员异常"))).collectList();
+    /**
+     *同步裂变计划
+     */
+    @Scheduled(fixedDelay = 110000)
+    @SchedulerLock(name = "launch:SyncPlanNewListenToRedis:syncGroupUsers", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
+    public void syncGroupUsers(){
+        Mono<List<GroupUsers>> listMono = template
+                .select(Query.query(Criteria.empty())
+                        .columns(COLUMNARRAY),GroupUsers.class)
+                .collectList();
+          listMono.flatMap(list -> {
+              Map<Integer, List<GroupUsers>> collect1 = list.stream().collect(Collectors.groupingBy(GroupUsers::groupUsersType));
+              collect1.forEach(((x,y)->
+                      {
+                          var groupUsersIdMysql = y.stream()
+                                  .map(GroupUsers::groupUsersId)
+                                  .map(String::valueOf).collect(Collectors.toList());
+                          reactiveRedisTemplate
+                                  .opsForSet()
+                                  .members(GroupUsersRedisKey.GROUP_USERS_INFO_SET+x)
+                                  .switchIfEmpty(Mono.just(""))
+                                  .collectList().flatMap(listIds -> {
+                              var addList = groupUsersIdMysql.stream().filter(f->!listIds.contains(f)).collect(Collectors.toList());
+                              if (!addList.isEmpty()) {
+                                  List<GroupUsers> collect = y.stream()
+                                          .filter(f -> addList.contains(String.valueOf(f.groupUsersId())))
+                                          .collect(Collectors.toList());
+                                  collect.forEach(groupUser-> setRedisValue(groupUser, String.valueOf(x)).subscribe());
+                              }
+                              var deleteList = listIds.stream().filter(f->!groupUsersIdMysql.contains(f)).collect(Collectors.toList());
+                              if (!deleteList.isEmpty()) {
+                                  deleteList.forEach(groupUser-> executeDelete(groupUser, String.valueOf(x)).subscribe());
+                              }
+                              return Mono.empty();
+                          }).subscribe();
+                      }
+                      ));
+              return Mono.empty();
+        }).subscribe();
+    }
+
+
+    @NotNull
+    public Mono<Void> setRedisValue(GroupUsers groupUsers,String groupUsersType) {
+        reactiveRedisTemplate.opsForSet().add(GroupUsersRedisKey.GROUP_USERS_INFO_SET+groupUsersType,String.valueOf( groupUsers.groupUsersId())).subscribe();
+        String json= getObjectToString(groupUsers);
+        reactiveRedisTemplate.opsForHash().put(GroupUsersRedisKey.GROUP_USERS_INFO_HASH,String.valueOf( groupUsers.groupUsersId()), json).subscribe();
+       return Mono.empty();
+    }
+    @NotNull
+    public Mono<Void> executeDelete(String groupUsersId,String groupUsersType) {
+        reactiveRedisTemplate.opsForSet().remove(GroupUsersRedisKey.GROUP_USERS_INFO_SET+groupUsersType, groupUsersId).subscribe();
+        reactiveRedisTemplate.opsForHash().remove(GroupUsersRedisKey.GROUP_USERS_INFO_HASH , groupUsersId).subscribe();
+        return Mono.empty();
+    }
+
+    private String getObjectToString(GroupUsers planNew) {
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
+        try {
+            // 将 record 转换为 JSON 字符串
+            return objectMapper.writeValueAsString(planNew);
+        } catch (JsonProcessingException e) {
+           throw new BaseException("PlanNew 将 record 转换为 JSON 字符串 异常:::"+e.getMessage());
+        }
+    }
+
+    /**
+     * 通过code从数据库获取PlanNew并缓存redis
+     * @param groupUsersType
+     * @return
+     */
+    public Mono<List<GroupUsers>> byCodeAndMysqlGetPlanNew(String groupUsersType,Integer limit){
+        Flux<GroupUsers> groupUsersFlux = template.select(Query.query(Criteria
+                .where("group_users_type").is(groupUsersType)).limit(limit), GroupUsers.class)
+                .onErrorResume(throwable -> Mono.error(new BaseException("获取群成员异常")));
+        return groupUsersFlux.flatMap(planNew -> {
+           setRedisValue(planNew,String.valueOf(planNew.groupUsersType())).subscribe();
+           return Mono.just(planNew);
+        }).collectList();
+    }
+
+    /**
+     * 通过redis获取缓存 空或失败 从mysql中取
+     *
+     * @param groupUsersType
+     * @param limit 截取长度
+     * @return {@link Mono}<{@link List}<{@link GroupUsers}>>
+     */
+    public Mono<List<GroupUsers>> byRedisGetGroupUsersIfErByMysql(String groupUsersType,Integer limit){
+      return reactiveRedisTemplate
+                .opsForSet()
+                .members(GroupUsersRedisKey.GROUP_USERS_INFO_SET+groupUsersType)
+                .collectList()
+                .switchIfEmpty(Mono.just(List.of()))
+              .onErrorResume(throwable -> {
+                  log.error("接口异常  {} e = {}","获取redis GroupUsers 失败 暂时从Mysql获取内容" , LogExceptionStackTrace.erroStackTrace(throwable));
+                  return Mono.just(List.of());
+              })
+              .flatMap(list -> {
+                     if( Objects.nonNull(list) && !list.isEmpty()){
+                       return   Flux.fromIterable(list).flatMap(f-> reactiveRedisTemplate.opsForHash().get(GroupUsersRedisKey.GROUP_USERS_INFO_HASH , f)).flatMap(fm->{
+                             GroupUsers planNew = strParsePlanNew(String.valueOf(fm));
+                             return Mono.just(planNew);
+                         }).collectList();
+                     }else {
+                         return byCodeAndMysqlGetPlanNew(groupUsersType,limit);
+                     }
+         });
+    }
+    /**
+     * str 转PlanNew
+     * @param str
+     * @return
+     */
+    public GroupUsers strParsePlanNew(String str) {
+        // 创建 ObjectMapper 实例并配置
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        SimpleModule module = new SimpleModule();
+        module.addDeserializer(ZonedDateTime.class, new CustomZonedDateTimeDeserializer());
+        objectMapper.registerModule(module);
+            try {
+                // 将 JSON 字符串转换为 record
+                return objectMapper.readValue(str, GroupUsers.class);
+            } catch (JsonProcessingException e) {
+                throw new BaseException(" 将 JSON字符串 转换为 GroupUsers  异常:::"+e.getMessage());
+            }
+    }
+    // 自定义反序列化器
+    public static class CustomZonedDateTimeDeserializer extends JsonDeserializer<ZonedDateTime> {
+        private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+        @Override
+        public ZonedDateTime deserialize(JsonParser p, DeserializationContext ctxt)
+                throws IOException, JsonProcessingException {
+            String date = p.getText();
+            LocalDateTime localDateTime = LocalDateTime.parse(date, FORMATTER);
+            // 使用系统默认时区,可以根据需要修改
+            return ZonedDateTime.of(localDateTime, ZoneId.systemDefault());
+        }
+    }
+}

+ 209 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncImageToRedis.java

@@ -0,0 +1,209 @@
+package com.webflux.launchadmin.mysql.task;
+
+import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationContext;
+import com.fasterxml.jackson.databind.JsonDeserializer;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import com.fasterxml.jackson.databind.module.SimpleModule;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
+import com.plumelog.core.util.LogExceptionStackTrace;
+import com.webflux.launchadmin.global.BaseException;
+import com.webflux.launchadmin.mysql.entity.listening.ListeningPlanNew;
+import com.webflux.launchadmin.mysql.entity.planNew.Image;
+import com.weblux.launchredis.utils.FissionRedisKey;
+import com.weblux.launchredis.utils.ImageRedisKey;
+import com.weblux.launchredis.utils.LinteningAudioRedisKey;
+import com.weblux.launchredis.utils.LinteningRedisKey;
+import jakarta.annotation.Resource;
+import lombok.extern.slf4j.Slf4j;
+import net.javacrumbs.shedlock.spring.annotation.SchedulerLock;
+import org.jetbrains.annotations.NotNull;
+import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
+import org.springframework.data.redis.core.ReactiveRedisTemplate;
+import org.springframework.data.relational.core.query.Criteria;
+import org.springframework.data.relational.core.query.Query;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+import java.io.IOException;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+@Component
+@Slf4j
+public class SyncImageToRedis {
+    @Resource
+    private R2dbcEntityTemplate template;
+    @Resource
+    private ReactiveRedisTemplate<String,String> reactiveRedisTemplate;
+
+
+    private static final String COLUMN="id,image,type,group_type,group_type_name,group_type_id,size_type";
+     private static final String[] COLUMNARRAY=COLUMN.split(",");
+      /**
+     *同步裂变计划
+     */
+//      Mono<List<Image>> imageMono = template.select(Query.query(Criteria.where("type").is(1)
+//              .and("group_type_id").is(f.groupId())
+//              .and("deleted_at").isNull()), Image.class)
+//              .onErrorResume(throwable -> Mono.error(new BaseException("获取群封面异常"))).collectList();
+      @Scheduled(fixedDelay = 120000)
+    @SchedulerLock(name = "launch:SyncImageToRedis:syncImagePlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
+    public void syncImagePlan(){
+        Query query = Query.query(Criteria
+                .where("deleted_at").isNull()
+                .and("group_type_id").isNotNull()
+                .and("type").isNotNull())
+                .columns(COLUMNARRAY);
+        Mono<List<Image>> listMono = template
+                .select(query,Image.class)
+                .collectList();
+        listMono.flatMap(list -> {
+            //根据groupId分组 在根据type分组
+            Map<Integer, List<Image>> collect1 = list.stream().collect(Collectors.groupingBy(Image::groupTypeId));
+            collect1.forEach(((x,y)->{
+                Map<Integer, List<Image>> collect = y.stream().collect(Collectors.groupingBy(Image::type));
+                collect.forEach((x1,y1)-> y1.forEach(image-> setRedisValue(image, x+"-"+x1).subscribe())); }));
+            return Mono.empty();
+        }).subscribe();
+        Query queryDelete = Query.query(Criteria
+                .where("deleted_at").isNotNull()
+                .and("group_type_id").isNotNull()
+                .and("type").isNotNull())
+                .columns(COLUMNARRAY);
+        Mono<List<Image>> listMonoDe = template
+                .select(queryDelete,Image.class)
+                .collectList();
+        listMonoDe.flatMap(list -> {
+            //根据groupId分组 在根据type分组
+            Map<Integer, List<Image>> collect1 = list.stream().collect(Collectors.groupingBy(Image::groupTypeId));
+            collect1.forEach(((x,y)->{
+                Map<Integer, List<Image>> collect = y.stream().collect(Collectors.groupingBy(Image::type));
+                collect.forEach((x1,y1)-> y1.forEach(image->
+                        deleteRedisValue(image, x+"-"+x1).subscribe()
+                ));
+            }));
+            return Mono.empty();
+        }).subscribe();
+    }
+
+
+    public Mono<Void> deleteRedisValue(Image image,String key) {
+        reactiveRedisTemplate.opsForSet().remove(ImageRedisKey.IMAGE_INFO_SET+key, String.valueOf(image.id())).subscribe();
+        reactiveRedisTemplate.opsForHash().remove(ImageRedisKey.IMAGE_INFO_HASH,   String.valueOf(image.id())).subscribe();
+        return Mono.empty();
+    }
+
+    public Mono<Void> setRedisValue(Image image,String key) {
+        reactiveRedisTemplate.opsForSet().add(ImageRedisKey.IMAGE_INFO_SET+key,String.valueOf(image.id()) ).subscribe();
+        String json= getImageToString(image);
+        reactiveRedisTemplate.opsForHash().put(ImageRedisKey.IMAGE_INFO_HASH,  String.valueOf(image.id()), json).subscribe();
+        return Mono.empty();
+    }
+
+
+
+    /**
+     * 通过code从数据库获取Image并缓存redis
+     * @return
+     */
+    public Mono<List<Image>> byCodeAndMysqlGetImage(String groupId,String type){
+        Mono<List<Image>> imageMono = template.select(Query.query(Criteria.where("type").is(type)
+                .and("group_type_id").is(groupId)
+                .and("deleted_at").isNull()), Image.class)
+                .onErrorResume(throwable -> Mono.error(new BaseException("获取群封面异常"))).collectList();
+       return imageMono.flatMap(images -> {
+           images.forEach(image -> {
+               setRedisValue(image,image.groupTypeId()+"-"+image.type()).subscribe();
+           });
+           return Mono.just(images);
+        });
+    }
+
+    /**
+     * 通过redis获取缓存 空或失败 从mysql中取
+     * @param groupId
+     * @param type
+     * @return {@link Mono}<{@link List}<{@link Image}>>
+     */
+    public Mono<List<Image>> byRedisGetImageIfErByMysql(String groupId,String type){
+        String key = groupId+"-"+type;
+        return reactiveRedisTemplate.opsForSet()
+                .members(ImageRedisKey.IMAGE_INFO_SET + key)
+                .collectList()
+                .switchIfEmpty(Mono.just(List.of()))
+                .onErrorResume(throwable -> {
+                    log.error("接口异常  {} e = {}", "获取redis Image 失败 暂时从Mysql获取内容", LogExceptionStackTrace.erroStackTrace(throwable));
+                    return Mono.just(List.of());
+                }).flatMap(listKey -> {
+                    if (Objects.nonNull(listKey) && !listKey.isEmpty()) {
+                        return Flux.fromIterable(listKey).flatMap(f -> reactiveRedisTemplate
+                                .opsForHash()
+                                .get(ImageRedisKey.IMAGE_INFO_HASH, f))
+                                .flatMap(str -> {
+                                    if (Objects.nonNull(str) && !"".equals(str)) {
+                                        return Mono.just(strParseImage(String.valueOf(str)));
+                                    }
+                                    return Mono.empty();
+                                }).collectList();
+                    } else {
+                        return byCodeAndMysqlGetImage(groupId, type);
+                    }
+                });
+    }
+
+
+    private String getImageToString(Image planNew) {
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
+        try {
+            // 将 record 转换为 JSON 字符串
+            return objectMapper.writeValueAsString(planNew);
+        } catch (JsonProcessingException e) {
+            throw new BaseException("Image 将 record 转换为 JSON 字符串 异常:::"+e.getMessage());
+        }
+    }
+    /**
+     * str 转Image
+     * @param str
+     * @return
+     */
+    public Image strParseImage(String str) {
+        // 创建 ObjectMapper 实例并配置
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        SimpleModule module = new SimpleModule();
+        module.addDeserializer(ZonedDateTime.class, new CustomZonedDateTimeDeserializer());
+        objectMapper.registerModule(module);
+            try {
+                // 将 JSON 字符串转换为 record
+                return objectMapper.readValue(str, Image.class);
+            } catch (JsonProcessingException e) {
+                throw new BaseException(" 将 JSON字符串 转换为 Image  异常:::"+e.getMessage());
+            }
+    }
+    // 自定义反序列化器
+    public static class CustomZonedDateTimeDeserializer extends JsonDeserializer<ZonedDateTime> {
+        private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+        @Override
+        public ZonedDateTime deserialize(JsonParser p, DeserializationContext ctxt)
+                throws IOException, JsonProcessingException {
+            String date = p.getText();
+            LocalDateTime localDateTime = LocalDateTime.parse(date, FORMATTER);
+            // 使用系统默认时区,可以根据需要修改
+            return ZonedDateTime.of(localDateTime, ZoneId.systemDefault());
+        }
+    }
+}

+ 1 - 1
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncListeningAudioToRedis.java

@@ -53,7 +53,7 @@ public class SyncListeningAudioToRedis {
     /**
      *同步裂变计划
      */
-    @Scheduled(cron="0 */30 * * * ?")
+    @Scheduled(fixedDelay = 130000)
     @SchedulerLock(name = "launch:SyncListeningAudioToRedis:syncLinteningAudioPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
     public void syncLinteningAudioPlan(){
         Mono<List<ListeningAudio>> listMono = template.select(Query.query(Criteria.empty())

+ 202 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncMaterialToRedis.java

@@ -0,0 +1,202 @@
+package com.webflux.launchadmin.mysql.task;
+
+import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationContext;
+import com.fasterxml.jackson.databind.JsonDeserializer;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import com.fasterxml.jackson.databind.module.SimpleModule;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
+import com.plumelog.core.util.LogExceptionStackTrace;
+import com.webflux.launchadmin.global.BaseException;
+import com.webflux.launchadmin.mysql.entity.planNew.Material;
+import com.weblux.launchredis.utils.MaterialRedisKey;
+import jakarta.annotation.Resource;
+import lombok.extern.slf4j.Slf4j;
+import net.javacrumbs.shedlock.spring.annotation.SchedulerLock;
+import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
+import org.springframework.data.redis.core.ReactiveRedisTemplate;
+import org.springframework.data.relational.core.query.Criteria;
+import org.springframework.data.relational.core.query.Query;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+import java.io.IOException;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+@Component
+@Slf4j
+public class SyncMaterialToRedis {
+    @Resource
+    private R2dbcEntityTemplate template;
+    @Resource
+    private ReactiveRedisTemplate<String,String> reactiveRedisTemplate;
+
+
+    private static final String COLUMN="id,content,type,group_type,group_type_name,group_type_id,remark";
+     private static final String[] COLUMNARRAY=COLUMN.split(",");
+      /**
+     *同步裂变计划
+     */
+//      Mono<List<Material>> textMono = template.select(Query.query(Criteria.where("type").is(4)
+//              .and("deleted_at").isNull()
+//              .and("group_type_id").is(f.groupId())), Material.class)
+      @Scheduled(fixedDelay = 140000)
+    @SchedulerLock(name = "launch:SyncMaterialToRedis:syncMaterialPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
+    public void syncMaterialPlan(){
+        Query query = Query.query(Criteria
+                .where("deleted_at").isNull()
+                .and("group_type_id").isNotNull()
+                .and("type").isNotNull())
+                .columns(COLUMNARRAY);
+        Mono<List<Material>> listMono = template
+                .select(query,Material.class)
+                .collectList();
+        listMono.flatMap(list -> {
+            //根据groupId分组 在根据type分组
+            Map<Integer, List<Material>> collect1 = list.stream().collect(Collectors.groupingBy(Material::groupTypeId));
+            collect1.forEach(((x,y)->{
+                Map<Integer, List<Material>> collect = y.stream().collect(Collectors.groupingBy(Material::type));
+                collect.forEach((x1,y1)-> y1.forEach(image-> setRedisValue(image, x+"-"+x1).subscribe())); }));
+            return Mono.empty();
+        }).subscribe();
+        Query queryDelete = Query.query(Criteria
+                .where("deleted_at").isNotNull()
+                .and("group_type_id").isNotNull()
+                .and("type").isNotNull())
+                .columns(COLUMNARRAY);
+        Mono<List<Material>> listMonoDe = template
+                .select(queryDelete,Material.class)
+                .collectList();
+        listMonoDe.flatMap(list -> {
+            //根据groupId分组 在根据type分组
+            Map<Integer, List<Material>> collect1 = list.stream().collect(Collectors.groupingBy(Material::groupTypeId));
+            collect1.forEach(((x,y)->{
+                Map<Integer, List<Material>> collect = y.stream().collect(Collectors.groupingBy(Material::type));
+                collect.forEach((x1,y1)-> y1.forEach(image->
+                        deleteRedisValue(image, x+"-"+x1).subscribe()
+                ));
+            }));
+            return Mono.empty();
+        }).subscribe();
+    }
+
+
+    public Mono<Void> deleteRedisValue(Material image,String key) {
+        reactiveRedisTemplate.opsForSet().remove(MaterialRedisKey.MATERIAL_INFO_SET+key, String.valueOf(image.id())).subscribe();
+        reactiveRedisTemplate.opsForHash().remove(MaterialRedisKey.MATERIAL_INFO_HASH,   String.valueOf(image.id())).subscribe();
+        return Mono.empty();
+    }
+
+    public Mono<Void> setRedisValue(Material image,String key) {
+        reactiveRedisTemplate.opsForSet().add(MaterialRedisKey.MATERIAL_INFO_SET+key,String.valueOf(image.id()) ).subscribe();
+        String json= getMaterialToString(image);
+        reactiveRedisTemplate.opsForHash().put(MaterialRedisKey.MATERIAL_INFO_HASH,  String.valueOf(image.id()), json).subscribe();
+        return Mono.empty();
+    }
+
+
+
+    /**
+     * 通过code从数据库获取Material并缓存redis
+     * @return
+     */
+    public Mono<List<Material>> byCodeAndMysqlGetMaterial(String groupId,String type){
+        Mono<List<Material>> imageMono = template.select(Query.query(Criteria.where("type").is(type)
+                .and("group_type_id").is(groupId)
+                .and("deleted_at").isNull()), Material.class)
+                .onErrorResume(throwable -> Mono.error(new BaseException("获取Material异常"))).collectList();
+       return imageMono.flatMap(images -> {
+           images.forEach(image -> {
+               setRedisValue(image,image.groupTypeId()+"-"+image.type()).subscribe();
+           });
+           return Mono.just(images);
+        });
+    }
+
+    /**
+     * 通过redis获取缓存 空或失败 从mysql中取
+     * @param groupId
+     * @param type
+     * @return {@link Mono}<{@link List}<{@link Material}>>
+     */
+    public Mono<List<Material>> byRedisGetMaterialIfErByMysql(String groupId,String type){
+        String key = groupId+"-"+type;
+        return reactiveRedisTemplate.opsForSet()
+                .members(MaterialRedisKey.MATERIAL_INFO_SET + key)
+                .collectList()
+                .switchIfEmpty(Mono.just(List.of()))
+                .onErrorResume(throwable -> {
+                    log.error("接口异常  {} e = {}", "获取redis Material 失败 暂时从Mysql获取内容", LogExceptionStackTrace.erroStackTrace(throwable));
+                    return Mono.just(List.of());
+                }).flatMap(listKey -> {
+                    if (Objects.nonNull(listKey) && !listKey.isEmpty()) {
+                        return Flux.fromIterable(listKey).flatMap(f -> reactiveRedisTemplate
+                                .opsForHash()
+                                .get(MaterialRedisKey.MATERIAL_INFO_HASH, f))
+                                .flatMap(str -> {
+                                    if (Objects.nonNull(str) && !"".equals(str)) {
+                                        return Mono.just(strParseMaterial(String.valueOf(str)));
+                                    }
+                                    return Mono.empty();
+                                }).collectList();
+                    } else {
+                        return byCodeAndMysqlGetMaterial(groupId, type);
+                    }
+                });
+    }
+
+
+    private String getMaterialToString(Material planNew) {
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
+        try {
+            // 将 record 转换为 JSON 字符串
+            return objectMapper.writeValueAsString(planNew);
+        } catch (JsonProcessingException e) {
+            throw new BaseException("Material 将 record 转换为 JSON 字符串 异常:::"+e.getMessage());
+        }
+    }
+    /**
+     * str 转Material
+     * @param str
+     * @return
+     */
+    public Material strParseMaterial(String str) {
+        // 创建 ObjectMapper 实例并配置
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        SimpleModule module = new SimpleModule();
+        module.addDeserializer(ZonedDateTime.class, new CustomZonedDateTimeDeserializer());
+        objectMapper.registerModule(module);
+            try {
+                // 将 JSON 字符串转换为 record
+                return objectMapper.readValue(str, Material.class);
+            } catch (JsonProcessingException e) {
+                throw new BaseException(" 将 JSON字符串 转换为 Material  异常:::"+e.getMessage());
+            }
+    }
+    // 自定义反序列化器
+    public static class CustomZonedDateTimeDeserializer extends JsonDeserializer<ZonedDateTime> {
+        private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+        @Override
+        public ZonedDateTime deserialize(JsonParser p, DeserializationContext ctxt)
+                throws IOException, JsonProcessingException {
+            String date = p.getText();
+            LocalDateTime localDateTime = LocalDateTime.parse(date, FORMATTER);
+            // 使用系统默认时区,可以根据需要修改
+            return ZonedDateTime.of(localDateTime, ZoneId.systemDefault());
+        }
+    }
+}

+ 3 - 3
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncPlanNewListenToRedis.java

@@ -51,7 +51,7 @@ public class SyncPlanNewListenToRedis {
     /**
      *同步裂变计划
      */
-    @Scheduled(cron="0 */29 * * * ?")
+    @Scheduled(fixedDelay = 160000)
     @SchedulerLock(name = "launch:SyncPlanNewListenToRedis:syncLinteningPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
     public void syncLinteningPlan(){
         System.out.println("start");
@@ -132,8 +132,8 @@ public class SyncPlanNewListenToRedis {
                 .switchIfEmpty(Mono.error(new BaseException("计划code异常 返回结果为空"))) //Mono.error(new BaseException("查询海报模版 id:"+f.posterTemplateId()+"=>空"))
                 .onErrorResume(throwable -> Mono.error(new BaseException(throwable.getMessage())));
        return planNewMono.flatMap(planNew -> {
-           Mono<Void> voidMono = setRedisValue(planNew);
-           return voidMono.then(Mono.just(planNew));
+            setRedisValue(planNew).subscribe();
+           return Mono.just(planNew);
         });
     }
 

+ 4 - 4
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncPlanNewToRedis.java

@@ -56,8 +56,8 @@ public class SyncPlanNewToRedis {
     /**
      *同步裂变计划
      */
-    //@Scheduled(cron="0 */1 * * * ?")
-   // @SchedulerLock(name = "launch:SyncPlanNewToRedis:syncFissionPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
+    @Scheduled(fixedDelay = 170000)
+    @SchedulerLock(name = "launch:SyncPlanNewToRedis:syncFissionPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
     public void syncFissionPlan(){
         Mono<List<PlanNew>> listMono = template.select(PlanNew.class).all().collectList();
         listMono.flatMap(list -> {
@@ -131,8 +131,8 @@ public class SyncPlanNewToRedis {
                 .switchIfEmpty(Mono.error(new BaseException("计划code异常 返回结果为空"))) //Mono.error(new BaseException("查询海报模版 id:"+f.posterTemplateId()+"=>空"))
                 .onErrorResume(throwable -> Mono.error(new BaseException(throwable.getMessage())));
        return planNewMono.flatMap(planNew -> {
-           Mono<Void> voidMono = setRedisValue(planNew);
-           return voidMono.then(Mono.just(planNew));
+           setRedisValue(planNew).subscribe();
+           return Mono.just(planNew);
         });
     }
 

+ 214 - 0
launch-admin/src/main/java/com/webflux/launchadmin/mysql/task/SyncReactiveDomainToRedis.java

@@ -0,0 +1,214 @@
+package com.webflux.launchadmin.mysql.task;
+
+import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationContext;
+import com.fasterxml.jackson.databind.JsonDeserializer;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import com.fasterxml.jackson.databind.module.SimpleModule;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
+import com.plumelog.core.util.LogExceptionStackTrace;
+import com.webflux.launchadmin.global.BaseException;
+import com.webflux.launchadmin.mysql.entity.planNew.DomainSelect;
+import com.weblux.launchredis.utils.ReactiveDomainRedisKey;
+import jakarta.annotation.Resource;
+import lombok.extern.slf4j.Slf4j;
+import net.javacrumbs.shedlock.spring.annotation.SchedulerLock;
+import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
+import org.springframework.data.redis.core.ReactiveRedisTemplate;
+import org.springframework.data.relational.core.query.Criteria;
+import org.springframework.data.relational.core.query.Query;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+import java.io.IOException;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+@Component
+@Slf4j
+public class SyncReactiveDomainToRedis {
+    @Resource
+    private R2dbcEntityTemplate template;
+    @Resource
+    private ReactiveRedisTemplate<String,String> reactiveRedisTemplate;
+
+
+    private static final String COLUMN="id,domain,type,status,group_type,group_type_id";
+     private static final String[] COLUMNARRAY=COLUMN.split(",");
+      /**
+     *同步裂变计划
+     */
+
+//    stringMono = template.select(Query.query(Criteria
+//            .where("group_type_id").is(f.transferPageDomain())
+//            .and("status").is(1)
+//                                .and("type").is(5)
+//                                .and("deleted_at").isNull()), DomainSelect.class)
+    @Scheduled(fixedDelay = 100000)
+    @SchedulerLock(name = "launch:SyncReactiveDomainToRedis:syncDomainSelectPlan", lockAtLeastFor = "PT5S", lockAtMostFor = "PT15M")
+    public void syncDomainSelectPlan(){
+        // group_type_id-type-status
+        Query query = Query.query(Criteria
+                .where("deleted_at").isNull()
+                .and("group_type_id").isNotNull()
+                .and("status").is(1)
+                .and("type").isNotNull())
+                .columns(COLUMNARRAY);
+        Mono<List<DomainSelect>> listMono = template
+                .select(query,DomainSelect.class)
+                .collectList();
+        listMono.flatMap(list -> {
+            //根据groupId分组 在根据type分组
+            var collect1 = list.stream().collect(Collectors.groupingBy(DomainSelect::groupTypeId));
+            collect1.forEach((x, y)->{
+                var collect = y.stream().collect(Collectors.groupingBy(DomainSelect::type));
+                collect.forEach((x1,y1)->{
+                    Map<String, List<DomainSelect>> collect2 = y1.stream().collect(Collectors.groupingBy(DomainSelect::status));
+                    collect2.forEach((x2,y2)->{
+                        y2.forEach(image-> setRedisValue(image, x+"-"+x1+"-"+x2).subscribe());
+                    }); }
+                ); });
+            return Mono.empty();
+        }).subscribe();
+        Query queryDelete = Query.query(Criteria
+                .where("deleted_at").isNotNull()
+                .and("group_type_id").isNotNull()
+                .and("status").is(1)
+                .and("type").isNotNull())
+                .columns(COLUMNARRAY);
+        Mono<List<DomainSelect>> listMonoDe = template
+                .select(queryDelete,DomainSelect.class)
+                .collectList();
+        listMonoDe.flatMap(list -> {
+            list.forEach(x->{
+                deleteRedisValue(x, x.groupTypeId()+"-"+x.groupType()+"-"+x.status()).subscribe();
+            });
+            return Mono.empty();
+        }).subscribe();
+
+    }
+
+
+    public Mono<Void> deleteRedisValue(DomainSelect image,String key) {
+        reactiveRedisTemplate.opsForSet().remove(ReactiveDomainRedisKey.DOMAIN_INFO_SET+key, String.valueOf(image.id())).subscribe();
+        reactiveRedisTemplate.opsForHash().remove(ReactiveDomainRedisKey.DOMAIN_INFO_HASH,   String.valueOf(image.id())).subscribe();
+        return Mono.empty();
+    }
+
+    public Mono<Void> setRedisValue(DomainSelect image,String key) {
+        reactiveRedisTemplate.opsForSet().add(ReactiveDomainRedisKey.DOMAIN_INFO_SET+key,String.valueOf(image.id()) ).subscribe();
+        String json= getDomainSelectToString(image);
+        reactiveRedisTemplate.opsForHash().put(ReactiveDomainRedisKey.DOMAIN_INFO_HASH,  String.valueOf(image.id()), json).subscribe();
+        return Mono.empty();
+    }
+
+
+
+    /**
+     * 通过code从数据库获取DomainSelect并缓存redis
+     * @return
+     */
+    public Mono<List<DomainSelect>> byCodeAndMysqlGetDomainSelect(String groupTypeId,String type ){
+        Mono<List<DomainSelect>> imageMono = template.select(Query.query(Criteria
+            .where("group_type_id").is(groupTypeId)
+            .and("status").is(1)
+            .and("type").is(type)
+            .and("deleted_at").isNull()), DomainSelect.class)
+                .onErrorResume(throwable -> Mono.error(new BaseException("byCodeAndMysqlGetDomainSelect 获取域名异常"+throwable.getMessage())))
+                .collectList();
+       return imageMono.flatMap(domains -> {
+           domains.forEach(domain -> {
+               setRedisValue(domain,domain.groupTypeId()+"-"+domain.type()+"-"+domain.status()).subscribe();
+           });
+           return Mono.just(domains);
+        });
+    }
+
+    /**
+     * 通过redis获取缓存 空或失败 从mysql中取
+     *
+     * @param groupTypeId
+     * @param type
+
+     * @return {@link Mono}<{@link List}<{@link DomainSelect}>>
+     */
+    public Mono<List<DomainSelect>> byRedisGetBackgroupDomainSelectIemIfErByMysql(String groupTypeId,String type){
+        String key = groupTypeId+"-"+type+"-1";
+        return reactiveRedisTemplate.opsForSet()
+                .members(ReactiveDomainRedisKey.DOMAIN_INFO_SET + key)
+                .collectList()
+                .switchIfEmpty(Mono.just(List.of()))
+                .onErrorResume(throwable -> {
+                    log.error("接口异常  {} e = {}", "获取redis DomainSelect 失败 暂时从Mysql获取内容", LogExceptionStackTrace.erroStackTrace(throwable));
+                    return Mono.just(List.of());
+                }).flatMap(listKey -> {
+                    if (Objects.nonNull(listKey) && !listKey.isEmpty()) {
+                        return Flux.fromIterable(listKey).flatMap(f -> reactiveRedisTemplate
+                                .opsForHash()
+                                .get(ReactiveDomainRedisKey.DOMAIN_INFO_HASH, f))
+                                .flatMap(str -> {
+                                    if (Objects.nonNull(str) && !"".equals(str)) {
+                                        return Mono.just(strParseDomainSelect(String.valueOf(str)));
+                                    }
+                                    return Mono.empty();
+                                }).collectList();
+                    } else {
+                        return byCodeAndMysqlGetDomainSelect(groupTypeId, type);
+                    }
+                });
+    }
+
+
+    private String getDomainSelectToString(DomainSelect planNew) {
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
+        try {
+            // 将 record 转换为 JSON 字符串
+            return objectMapper.writeValueAsString(planNew);
+        } catch (JsonProcessingException e) {
+            throw new BaseException("DomainSelect 将 record 转换为 JSON 字符串 异常:::"+e.getMessage());
+        }
+    }
+    /**
+     * str 转DomainSelect
+     * @param str
+     * @return
+     */
+    public DomainSelect strParseDomainSelect(String str) {
+        // 创建 ObjectMapper 实例并配置
+        ObjectMapper objectMapper = new ObjectMapper();
+        objectMapper.registerModule(new JavaTimeModule());
+        SimpleModule module = new SimpleModule();
+        module.addDeserializer(ZonedDateTime.class, new CustomZonedDateTimeDeserializer());
+        objectMapper.registerModule(module);
+            try {
+                // 将 JSON 字符串转换为 record
+                return objectMapper.readValue(str, DomainSelect.class);
+            } catch (JsonProcessingException e) {
+                throw new BaseException(" 将 JSON字符串 转换为 DomainSelect  异常:::"+e.getMessage());
+            }
+    }
+    // 自定义反序列化器
+    public static class CustomZonedDateTimeDeserializer extends JsonDeserializer<ZonedDateTime> {
+        private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+        @Override
+        public ZonedDateTime deserialize(JsonParser p, DeserializationContext ctxt)
+                throws IOException, JsonProcessingException {
+            String date = p.getText();
+            LocalDateTime localDateTime = LocalDateTime.parse(date, FORMATTER);
+            // 使用系统默认时区,可以根据需要修改
+            return ZonedDateTime.of(localDateTime, ZoneId.systemDefault());
+        }
+    }
+}

+ 1 - 1
launch-admin/src/test/java/com/webflux/launchadmin/LaunchAdminApplicationTests.java

@@ -1,10 +1,10 @@
 package com.webflux.launchadmin;
 
+
 import com.plumelog.core.util.LogExceptionStackTrace;
 import org.junit.jupiter.api.Test;
 import org.springframework.boot.test.context.SpringBootTest;
 import reactor.core.publisher.Mono;
-
 import java.net.URI;
 import java.net.http.HttpClient;
 import java.net.http.HttpHeaders;

+ 15 - 0
launch-redis/src/main/java/com/weblux/launchredis/utils/BackgroupImageItemRedisKey.java

@@ -0,0 +1,15 @@
+package com.weblux.launchredis.utils;
+
+/**
+ * 裂变redis-key
+ *
+ * @author:
+ * @date: 2024/7/12 10:06
+ */
+public class BackgroupImageItemRedisKey {
+
+    public final static String BACKGROUP_IMAGE_ITEM_INFO = "sync_backgroup_image_item_info";
+    public final static String BACKGROUP_IMAGE_ITEM_INFO_SET = "sync_backgroup_image_item_set_ids_";
+    public final static String BACKGROUP_IMAGE_ITEM_INFO_KEY = "sync_backgroup_image_item_info_";
+    public final static String BACKGROUP_IMAGE_ITEM_INFO_HASH = "sync_backgroup_image_item_info_hash_json";
+}

+ 17 - 0
launch-redis/src/main/java/com/weblux/launchredis/utils/GroupUsersRedisKey.java

@@ -0,0 +1,17 @@
+package com.weblux.launchredis.utils;
+
+/**
+ * 裂变redis-key
+ *
+ * @author:
+ * @date: 2024/7/12 10:06
+ */
+public class GroupUsersRedisKey {
+
+
+    //裂变计划缓存
+    public final static String GROUP_USERS_INFO = "sync_group_users_info";
+    public final static String GROUP_USERS_INFO_SET = "sync_group_users_info_set_ids";
+    public final static String GROUP_USERS_INFO_KEY = "sync_group_users_info_";
+    public final static String GROUP_USERS_INFO_HASH = "sync_group_users_info_hash_json";
+}

+ 17 - 0
launch-redis/src/main/java/com/weblux/launchredis/utils/ImageRedisKey.java

@@ -0,0 +1,17 @@
+package com.weblux.launchredis.utils;
+
+/**
+ * 裂变redis-key
+ *
+ * @author:
+ * @date: 2024/7/12 10:06
+ */
+public class ImageRedisKey {
+
+
+    //裂变计划缓存
+    public final static String IMAGE_INFO = "sync_image_info";
+    public final static String IMAGE_INFO_SET = "sync_image_info_set_ids_";
+    public final static String IMAGE_INFO_HASH = "sync_image_info_hash_json";
+    public final static String IMAGE_INFO_KEY = "sync_image_info_id_";
+}

+ 17 - 0
launch-redis/src/main/java/com/weblux/launchredis/utils/MaterialRedisKey.java

@@ -0,0 +1,17 @@
+package com.weblux.launchredis.utils;
+
+/**
+ * 裂变redis-key
+ *
+ * @author:
+ * @date: 2024/7/12 10:06
+ */
+public class MaterialRedisKey {
+
+
+    //裂变计划缓存
+    public final static String MATERIAL_INFO = "sync_material_info";
+    public final static String MATERIAL_INFO_SET = "sync_material_info_set_ids_";
+    public final static String MATERIAL_INFO_HASH = "sync_material_info_hash_json";
+    public final static String MATERIAL_INFO_KEY = "sync_material_info_id_";
+}

+ 17 - 0
launch-redis/src/main/java/com/weblux/launchredis/utils/ReactiveDomainRedisKey.java

@@ -0,0 +1,17 @@
+package com.weblux.launchredis.utils;
+
+/**
+ * 裂变redis-key
+ *
+ * @author:
+ * @date: 2024/7/12 10:06
+ */
+public class ReactiveDomainRedisKey {
+
+
+    //裂变计划缓存
+    public final static String DOMAIN_INFO = "sync_reactive_domain_info";
+    public final static String DOMAIN_INFO_SET = "sync_reactive_domain_info_set_ids_";
+    public final static String DOMAIN_INFO_HASH = "sync_reactive_domain_info_hash_json";
+    public final static String DOMAIN_INFO_KEY = "sync_reactive_domain_info_id_";
+}