Explorar el Código

数据同步消息发送注解与逻辑

bing hace 1 año
padre
commit
65dec9383d

+ 19 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/aop/annotation/SyncDataToLocal.java

@@ -0,0 +1,19 @@
+package org.dromara.backstage.aop.annotation;
+
+import java.lang.annotation.*;
+
+/**
+ * 数据同步至本地服务数据库注解
+ *
+ * @author bing
+ */
+@Target({ElementType.METHOD})
+@Retention(RetentionPolicy.RUNTIME)
+@Documented
+public @interface SyncDataToLocal {
+    /**
+     * 消息的event_type
+     */
+    String eventType() default "";
+
+}

+ 96 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/aop/aspect/SyncDataToLocalAspect.java

@@ -0,0 +1,96 @@
+package org.dromara.backstage.aop.aspect;
+
+import cn.hutool.core.lang.UUID;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.StringUtils;
+import org.aspectj.lang.JoinPoint;
+import org.aspectj.lang.annotation.AfterReturning;
+import org.aspectj.lang.annotation.Aspect;
+import org.aspectj.lang.reflect.CodeSignature;
+import org.dromara.backstage.aop.annotation.SyncDataToLocal;
+import org.dromara.backstage.mq.KafkaProducer;
+import org.dromara.common.core.domain.R;
+import org.dromara.common.message.kafka.constant.KafkaTopicConstants;
+import org.dromara.common.message.kafka.domain.KafkaHeader;
+import org.dromara.common.message.kafka.domain.KafkaMessage;
+import org.dromara.common.satoken.utils.LoginHelper;
+import org.springframework.boot.autoconfigure.AutoConfiguration;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * 同步数据至本地 的 切面
+ *
+ * @author bing
+ */
+@RequiredArgsConstructor
+@Slf4j
+@Aspect
+@AutoConfiguration
+public class SyncDataToLocalAspect {
+
+    private final KafkaProducer kafkaProducer;
+
+
+    /**
+     * 处理完请求后执行
+     *
+     * @param joinPoint 切点
+     */
+    @AfterReturning(pointcut = "@annotation(controllerSyncData2Local)", returning = "jsonResult")
+    public void doAfterReturning(JoinPoint joinPoint, SyncDataToLocal controllerSyncData2Local, Object jsonResult) {
+        if (jsonResult instanceof R<?> r) {
+            if (r.getCode() == R.SUCCESS) {
+                sendSyncMessage(joinPoint, controllerSyncData2Local);
+            }else{
+                log.error("同步数据消息未发送:controller 方法返回结果未失败!");
+            }
+        }else{
+            log.error("同步数据消息未发送:controller 方法返回不是R类型!");
+        }
+
+    }
+
+    private void sendSyncMessage(JoinPoint joinPoint, SyncDataToLocal controllerSyncData2Local) {
+        try {
+            KafkaMessage<Object> data = new KafkaMessage<>();
+            KafkaHeader header = data.getHeader();
+            header.setTimestamp(System.currentTimeMillis());
+            header.setEventId(UUID.randomUUID().toString());
+            header.setEventType(controllerSyncData2Local.eventType());
+            String sender = header.getSender();
+            String eventType = header.getEventType();
+            if(StringUtils.isBlank(sender) && StringUtils.isNotBlank(eventType)){
+                header.setSender(eventType.substring(0, eventType.lastIndexOf("_")));
+            }
+            String tenantId = header.getTenantId();
+            if(StringUtils.isBlank(tenantId)){
+                header.setTenantId(LoginHelper.getTenantId());
+            }
+
+            Object[] args = joinPoint.getArgs();
+            int length = args.length;
+            if(length == 1){
+                data.setBody(args[0]);
+            }else if(length >1){
+                CodeSignature signature = (CodeSignature) joinPoint.getSignature();
+                String[] paramNames = signature.getParameterNames();
+                Map<String, Object> params = new HashMap<>();
+                for (int i = 0; i < length; i++) {
+//                System.out.println("参数名: " + paramNames[i] + ", 参数值: " + args[i]);
+                    params.put(paramNames[i], args[i]);
+                }
+                data.setBody(params);
+            }else{
+                data.setBody(null);
+            }
+            kafkaProducer.sendKafkaMessage(KafkaTopicConstants.SYNC_DATA_TOPIC, data);
+        }catch (Exception e){
+            log.error("同步数据消息未发送:发送消息异常!", e);
+        }
+    }
+
+
+}

+ 106 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/controller/SendMessageRecordController.java

@@ -0,0 +1,106 @@
+package org.dromara.backstage.basics.controller;
+
+import java.util.List;
+
+import lombok.RequiredArgsConstructor;
+import jakarta.servlet.http.HttpServletResponse;
+import jakarta.validation.constraints.*;
+import cn.dev33.satoken.annotation.SaCheckPermission;
+import org.springframework.web.bind.annotation.*;
+import org.springframework.validation.annotation.Validated;
+import org.dromara.common.idempotent.annotation.RepeatSubmit;
+import org.dromara.common.log.annotation.Log;
+import org.dromara.common.web.core.BaseController;
+import org.dromara.common.mybatis.core.page.PageQuery;
+import org.dromara.common.core.domain.R;
+import org.dromara.common.core.validate.AddGroup;
+import org.dromara.common.core.validate.EditGroup;
+import org.dromara.common.log.enums.BusinessType;
+import org.dromara.common.excel.utils.ExcelUtil;
+import org.dromara.backstage.basics.domain.vo.SendMessageRecordVo;
+import org.dromara.backstage.basics.domain.bo.SendMessageRecordBo;
+import org.dromara.backstage.basics.service.ISendMessageRecordService;
+import org.dromara.common.mybatis.core.page.TableDataInfo;
+
+/**
+ * 消息发送记录
+ * 前端访问路由地址为:/basics/sendMessageRecord
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+@Validated
+@RequiredArgsConstructor
+@RestController
+@RequestMapping("/basics/sendMessageRecord")
+public class SendMessageRecordController extends BaseController {
+
+    private final ISendMessageRecordService sendMessageRecordService;
+
+    /**
+     * 查询消息发送记录列表
+     */
+    @SaCheckPermission("basics:sendMessageRecord:list")
+    @GetMapping("/list")
+    public TableDataInfo<SendMessageRecordVo> list(SendMessageRecordBo bo, PageQuery pageQuery) {
+        return sendMessageRecordService.queryPageList(bo, pageQuery);
+    }
+
+    /**
+     * 导出消息发送记录列表
+     */
+    @SaCheckPermission("basics:sendMessageRecord:export")
+    @Log(title = "消息发送记录", businessType = BusinessType.EXPORT)
+    @PostMapping("/export")
+    public void export(SendMessageRecordBo bo, HttpServletResponse response) {
+        List<SendMessageRecordVo> list = sendMessageRecordService.queryList(bo);
+        ExcelUtil.exportExcel(list, "消息发送记录", SendMessageRecordVo.class, response);
+    }
+
+    /**
+     * 获取消息发送记录详细信息
+     *
+     * @param recordId 主键
+     */
+    @SaCheckPermission("basics:sendMessageRecord:query")
+    @GetMapping("/{recordId}")
+    public R<SendMessageRecordVo> getInfo(@NotNull(message = "主键不能为空")
+                                     @PathVariable Long recordId) {
+        return R.ok(sendMessageRecordService.queryById(recordId));
+    }
+
+    /**
+     * 新增消息发送记录
+     */
+    @SaCheckPermission("basics:sendMessageRecord:add")
+    @Log(title = "消息发送记录", businessType = BusinessType.INSERT)
+    @RepeatSubmit()
+    @PostMapping()
+    public R<Void> add(@Validated(AddGroup.class) @RequestBody SendMessageRecordBo bo) {
+        return toAjax(sendMessageRecordService.insertByBo(bo));
+    }
+
+    /**
+     * 修改消息发送记录
+     */
+    @SaCheckPermission("basics:sendMessageRecord:edit")
+    @Log(title = "消息发送记录", businessType = BusinessType.UPDATE)
+    @RepeatSubmit()
+    @PutMapping()
+    public R<Void> edit(@Validated(EditGroup.class) @RequestBody SendMessageRecordBo bo) {
+        return toAjax(sendMessageRecordService.updateByBo(bo));
+    }
+
+    /**
+     * 删除消息发送记录
+     *
+     * @param recordIds 主键串
+     */
+    @SaCheckPermission("basics:sendMessageRecord:remove")
+    @Log(title = "消息发送记录", businessType = BusinessType.DELETE)
+    @DeleteMapping("/{recordIds}")
+    public R<Void> remove(@NotEmpty(message = "主键不能为空")
+                          @PathVariable Long[] recordIds) {
+        return toAjax(sendMessageRecordService.deleteWithValidByIds(List.of(recordIds), true));
+    }
+}

+ 65 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/domain/SendMessageRecord.java

@@ -0,0 +1,65 @@
+package org.dromara.backstage.basics.domain;
+
+import org.dromara.common.tenant.core.TenantEntity;
+import com.baomidou.mybatisplus.annotation.*;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+
+import java.io.Serial;
+
+/**
+ * 消息发送记录对象 t_send_message_record
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@TableName("t_send_message_record")
+public class SendMessageRecord extends TenantEntity {
+
+    @Serial
+    private static final long serialVersionUID = 1L;
+
+    /**
+     * 主键id
+     */
+    private Long recordId;
+
+    /**
+     * 消息类型:kafka、rabbitmq、rocketmq
+     */
+    private String mqType;
+
+    /**
+     * 消息主题
+     */
+    private String topic;
+
+    /**
+     * 消息事件类型
+     */
+    private String eventType;
+
+    /**
+     * 发送结果:S 成功,F 失败
+     */
+    private String result;
+
+    /**
+     * 消息
+     */
+    private String message;
+
+    /**
+     * $column.columnComment
+     */
+    private String eventId;
+
+    /**
+     * 发送方
+     */
+    private String sender;
+
+
+}

+ 73 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/domain/bo/SendMessageRecordBo.java

@@ -0,0 +1,73 @@
+package org.dromara.backstage.basics.domain.bo;
+
+import org.dromara.backstage.basics.domain.SendMessageRecord;
+import org.dromara.common.mybatis.core.domain.BaseEntity;
+import org.dromara.common.core.validate.AddGroup;
+import org.dromara.common.core.validate.EditGroup;
+import io.github.linpeilie.annotations.AutoMapper;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import jakarta.validation.constraints.*;
+import org.dromara.common.tenant.core.TenantEntity;
+
+/**
+ * 消息发送记录业务对象 t_send_message_record
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@AutoMapper(target = SendMessageRecord.class, reverseConvertGenerate = false)
+public class SendMessageRecordBo extends TenantEntity {
+
+    /**
+     * 主键id
+     */
+    @NotNull(message = "主键id不能为空", groups = { AddGroup.class, EditGroup.class })
+    private Long recordId;
+
+    /**
+     * 消息类型:kafka、rabbitmq、rocketmq
+     */
+    @NotBlank(message = "消息类型:kafka、rabbitmq、rocketmq不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String mqType;
+
+    /**
+     * 消息主题
+     */
+    @NotBlank(message = "消息主题不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String topic;
+
+    /**
+     * 消息事件类型
+     */
+    @NotBlank(message = "消息事件类型不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String eventType;
+
+    /**
+     * 发送结果:S 成功,F 失败
+     */
+    @NotBlank(message = "发送结果:S 成功,F 失败不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String result;
+
+    /**
+     * 消息
+     */
+    @NotBlank(message = "消息不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String message;
+
+    /**
+     * $column.columnComment
+     */
+    @NotBlank(message = "$column.columnComment不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String eventId;
+
+    /**
+     * 发送方
+     */
+    @NotBlank(message = "发送方不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String sender;
+
+
+}

+ 81 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/domain/vo/SendMessageRecordVo.java

@@ -0,0 +1,81 @@
+package org.dromara.backstage.basics.domain.vo;
+
+import org.dromara.backstage.basics.domain.SendMessageRecord;
+import com.alibaba.excel.annotation.ExcelIgnoreUnannotated;
+import com.alibaba.excel.annotation.ExcelProperty;
+import org.dromara.common.excel.annotation.ExcelDictFormat;
+import org.dromara.common.excel.convert.ExcelDictConvert;
+import io.github.linpeilie.annotations.AutoMapper;
+import lombok.Data;
+
+import java.io.Serial;
+import java.io.Serializable;
+import java.util.Date;
+
+
+
+/**
+ * 消息发送记录视图对象 t_send_message_record
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+@Data
+@ExcelIgnoreUnannotated
+@AutoMapper(target = SendMessageRecord.class)
+public class SendMessageRecordVo implements Serializable {
+
+    @Serial
+    private static final long serialVersionUID = 1L;
+
+    /**
+     * 主键id
+     */
+    @ExcelProperty(value = "主键id")
+    private Long recordId;
+
+    /**
+     * 消息类型:kafka、rabbitmq、rocketmq
+     */
+    @ExcelProperty(value = "消息类型:kafka、rabbitmq、rocketmq")
+    private String mqType;
+
+    /**
+     * 消息主题
+     */
+    @ExcelProperty(value = "消息主题")
+    private String topic;
+
+    /**
+     * 消息事件类型
+     */
+    @ExcelProperty(value = "消息事件类型")
+    private String eventType;
+
+    /**
+     * 发送结果:S 成功,F 失败
+     */
+    @ExcelProperty(value = "发送结果:S 成功,F 失败")
+    private String result;
+
+    /**
+     * 消息
+     */
+    @ExcelProperty(value = "消息")
+    private String message;
+
+    /**
+     * $column.columnComment
+     */
+    @ExcelProperty(value = "${comment}", converter = ExcelDictConvert.class)
+    @ExcelDictFormat(readConverterExp = "$column.readConverterExp()")
+    private String eventId;
+
+    /**
+     * 发送方
+     */
+    @ExcelProperty(value = "发送方")
+    private String sender;
+
+
+}

+ 15 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/mapper/SendMessageRecordMapper.java

@@ -0,0 +1,15 @@
+package org.dromara.backstage.basics.mapper;
+
+import org.dromara.backstage.basics.domain.SendMessageRecord;
+import org.dromara.backstage.basics.domain.vo.SendMessageRecordVo;
+import org.dromara.common.mybatis.core.mapper.BaseMapperPlus;
+
+/**
+ * 消息发送记录Mapper接口
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+public interface SendMessageRecordMapper extends BaseMapperPlus<SendMessageRecord, SendMessageRecordVo> {
+
+}

+ 69 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/service/ISendMessageRecordService.java

@@ -0,0 +1,69 @@
+package org.dromara.backstage.basics.service;
+
+import org.dromara.backstage.basics.domain.SendMessageRecord;
+import org.dromara.backstage.basics.domain.vo.SendMessageRecordVo;
+import org.dromara.backstage.basics.domain.bo.SendMessageRecordBo;
+import org.dromara.common.mybatis.core.page.TableDataInfo;
+import org.dromara.common.mybatis.core.page.PageQuery;
+
+import java.util.Collection;
+import java.util.List;
+
+/**
+ * 消息发送记录Service接口
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+public interface ISendMessageRecordService {
+
+    /**
+     * 查询消息发送记录
+     *
+     * @param recordId 主键
+     * @return 消息发送记录
+     */
+    SendMessageRecordVo queryById(Long recordId);
+
+    /**
+     * 分页查询消息发送记录列表
+     *
+     * @param bo        查询条件
+     * @param pageQuery 分页参数
+     * @return 消息发送记录分页列表
+     */
+    TableDataInfo<SendMessageRecordVo> queryPageList(SendMessageRecordBo bo, PageQuery pageQuery);
+
+    /**
+     * 查询符合条件的消息发送记录列表
+     *
+     * @param bo 查询条件
+     * @return 消息发送记录列表
+     */
+    List<SendMessageRecordVo> queryList(SendMessageRecordBo bo);
+
+    /**
+     * 新增消息发送记录
+     *
+     * @param bo 消息发送记录
+     * @return 是否新增成功
+     */
+    Boolean insertByBo(SendMessageRecordBo bo);
+
+    /**
+     * 修改消息发送记录
+     *
+     * @param bo 消息发送记录
+     * @return 是否修改成功
+     */
+    Boolean updateByBo(SendMessageRecordBo bo);
+
+    /**
+     * 校验并批量删除消息发送记录信息
+     *
+     * @param ids     待删除的主键集合
+     * @param isValid 是否进行有效性校验
+     * @return 是否删除成功
+     */
+    Boolean deleteWithValidByIds(Collection<Long> ids, Boolean isValid);
+}

+ 154 - 0
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/basics/service/impl/SendMessageRecordServiceImpl.java

@@ -0,0 +1,154 @@
+package org.dromara.backstage.basics.service.impl;
+
+import org.dromara.common.core.utils.MapstructUtils;
+import org.dromara.common.core.utils.StringUtils;
+import org.dromara.common.mybatis.core.page.TableDataInfo;
+import org.dromara.common.mybatis.core.page.PageQuery;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import lombok.RequiredArgsConstructor;
+import org.springframework.stereotype.Service;
+import org.dromara.backstage.basics.domain.bo.SendMessageRecordBo;
+import org.dromara.backstage.basics.domain.vo.SendMessageRecordVo;
+import org.dromara.backstage.basics.domain.SendMessageRecord;
+import org.dromara.backstage.basics.mapper.SendMessageRecordMapper;
+import org.dromara.backstage.basics.service.ISendMessageRecordService;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Collection;
+
+/**
+ * 消息发送记录Service业务层处理
+ *
+ * @author bing
+ * @date 2024-10-30
+ */
+@RequiredArgsConstructor
+@Service
+public class SendMessageRecordServiceImpl implements ISendMessageRecordService {
+
+    private final SendMessageRecordMapper baseMapper;
+
+    /**
+     * 查询消息发送记录
+     *
+     * @param recordId 主键
+     * @return 消息发送记录
+     */
+    @Override
+    public SendMessageRecordVo queryById(Long recordId){
+        return baseMapper.selectVoById(recordId);
+    }
+
+    /**
+     * 分页查询消息发送记录列表
+     *
+     * @param bo        查询条件
+     * @param pageQuery 分页参数
+     * @return 消息发送记录分页列表
+     */
+    @Override
+    public TableDataInfo<SendMessageRecordVo> queryPageList(SendMessageRecordBo bo, PageQuery pageQuery) {
+        LambdaQueryWrapper<SendMessageRecord> lqw = buildQueryWrapper(bo);
+        Page<SendMessageRecordVo> result = baseMapper.selectVoPage(pageQuery.build(), lqw);
+        return TableDataInfo.build(result);
+    }
+
+    /**
+     * 查询符合条件的消息发送记录列表
+     *
+     * @param bo 查询条件
+     * @return 消息发送记录列表
+     */
+    @Override
+    public List<SendMessageRecordVo> queryList(SendMessageRecordBo bo) {
+        LambdaQueryWrapper<SendMessageRecord> lqw = buildQueryWrapper(bo);
+        return baseMapper.selectVoList(lqw);
+    }
+
+    private LambdaQueryWrapper<SendMessageRecord> buildQueryWrapper(SendMessageRecordBo bo) {
+        Map<String, Object> params = bo.getParams();
+        LambdaQueryWrapper<SendMessageRecord> lqw = Wrappers.lambdaQuery();
+        lqw.eq(bo.getRecordId() != null, SendMessageRecord::getRecordId, bo.getRecordId());
+        lqw.eq(StringUtils.isNotBlank(bo.getMqType()), SendMessageRecord::getMqType, bo.getMqType());
+        lqw.eq(StringUtils.isNotBlank(bo.getTopic()), SendMessageRecord::getTopic, bo.getTopic());
+        lqw.eq(StringUtils.isNotBlank(bo.getEventType()), SendMessageRecord::getEventType, bo.getEventType());
+        lqw.eq(StringUtils.isNotBlank(bo.getResult()), SendMessageRecord::getResult, bo.getResult());
+        lqw.eq(StringUtils.isNotBlank(bo.getMessage()), SendMessageRecord::getMessage, bo.getMessage());
+        lqw.eq(StringUtils.isNotBlank(bo.getEventId()), SendMessageRecord::getEventId, bo.getEventId());
+        lqw.eq(StringUtils.isNotBlank(bo.getSender()), SendMessageRecord::getSender, bo.getSender());
+        return lqw;
+    }
+
+    private QueryWrapper<SendMessageRecord> buildQueryWrapper(SendMessageRecordBo bo,String tableAlias) {
+        QueryWrapper<SendMessageRecord> lqw = new QueryWrapper<>();
+        String columnPrefix = "";
+        if(StringUtils.isNotBlank(tableAlias)){
+            columnPrefix = tableAlias + ".";
+        }
+        lqw.eq(bo.getRecordId() != null, columnPrefix+"record_id", bo.getRecordId());
+        lqw.eq(StringUtils.isNotBlank(bo.getMqType()), columnPrefix+"mq_type", bo.getMqType());
+        lqw.eq(StringUtils.isNotBlank(bo.getTopic()), columnPrefix+"topic", bo.getTopic());
+        lqw.eq(StringUtils.isNotBlank(bo.getEventType()), columnPrefix+"event_type", bo.getEventType());
+        lqw.eq(StringUtils.isNotBlank(bo.getResult()), columnPrefix+"result", bo.getResult());
+        lqw.eq(StringUtils.isNotBlank(bo.getMessage()), columnPrefix+"message", bo.getMessage());
+        lqw.eq(StringUtils.isNotBlank(bo.getEventId()), columnPrefix+"event_id", bo.getEventId());
+        lqw.eq(StringUtils.isNotBlank(bo.getSender()), columnPrefix+"sender", bo.getSender());
+        return lqw;
+    }
+
+    /**
+     * 新增消息发送记录
+     *
+     * @param bo 消息发送记录
+     * @return 是否新增成功
+     */
+    @Override
+    public Boolean insertByBo(SendMessageRecordBo bo) {
+        SendMessageRecord add = MapstructUtils.convert(bo, SendMessageRecord.class);
+        validEntityBeforeSave(add);
+        boolean flag = baseMapper.insert(add) > 0;
+        if (flag) {
+            bo.setRecordId(add.getRecordId());
+        }
+        return flag;
+    }
+
+    /**
+     * 修改消息发送记录
+     *
+     * @param bo 消息发送记录
+     * @return 是否修改成功
+     */
+    @Override
+    public Boolean updateByBo(SendMessageRecordBo bo) {
+        SendMessageRecord update = MapstructUtils.convert(bo, SendMessageRecord.class);
+        validEntityBeforeSave(update);
+        return baseMapper.updateById(update) > 0;
+    }
+
+    /**
+     * 保存前的数据校验
+     */
+    private void validEntityBeforeSave(SendMessageRecord entity){
+        //做一些数据校验,如唯一约束
+    }
+
+    /**
+     * 校验并批量删除消息发送记录信息
+     *
+     * @param ids     待删除的主键集合
+     * @param isValid 是否进行有效性校验
+     * @return 是否删除成功
+     */
+    @Override
+    public Boolean deleteWithValidByIds(Collection<Long> ids, Boolean isValid) {
+        if(isValid){
+            //做一些业务上的校验,判断是否需要校验
+        }
+        return baseMapper.deleteByIds(ids) > 0;
+    }
+}

+ 45 - 9
ruoyi-modules/ruoyi-backstage/src/main/java/org/dromara/backstage/mq/KafkaProducer.java

@@ -1,15 +1,23 @@
 package org.dromara.backstage.mq;
 
+import cn.hutool.core.lang.UUID;
+import com.alibaba.excel.util.StringUtils;
 import com.alibaba.fastjson2.JSON;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.kafka.clients.producer.ProducerRecord;
+import org.dromara.backstage.basics.domain.bo.SendMessageRecordBo;
+import org.dromara.backstage.basics.service.ISendMessageRecordService;
+import org.dromara.common.message.kafka.constant.KafkaTopicConstants;
 import org.dromara.common.message.kafka.domain.KafkaHeader;
 import org.dromara.common.message.kafka.domain.KafkaMessage;
+import org.dromara.common.satoken.utils.LoginHelper;
+import org.dromara.system.api.model.LoginUser;
 import org.springframework.kafka.core.KafkaTemplate;
 import org.springframework.kafka.support.SendResult;
 import org.springframework.stereotype.Component;
 
+import java.util.Date;
 import java.util.concurrent.CompletableFuture;
 
 @RequiredArgsConstructor
@@ -19,6 +27,8 @@ public class KafkaProducer {
 
     private final KafkaTemplate<String, String> kafkaTemplate;
 
+    private final ISendMessageRecordService sendMessageRecordService;
+
     /**
      * Send.
      *
@@ -31,13 +41,10 @@ public class KafkaProducer {
         log.debug("发送消息到kafka消息系统结束");
     }
 
-    public void sendSyncData(String topic, KafkaMessage<?> data){
+
+
+    public void sendKafkaMessage(String topic, KafkaMessage<?> data){
         try{
-            KafkaHeader header = data.getHeader();
-            String eventId = header.getEventId();
-            String sender = header.getSender();
-            String eventType = header.getEventType();
-            String tenantId = header.getTenantId();
             String jsonMessage = JSON.toJSONString(data);
             ProducerRecord<String, String> record = new ProducerRecord<>(topic, "YKT-SYNC-Message", jsonMessage);
             log.info("发送同步数据到kafka消息系统, data: " + jsonMessage);
@@ -46,17 +53,46 @@ public class KafkaProducer {
                 if (ex != null) {
                     log.error("同步数据发送到kafka消息系统异常,data: " + jsonMessage, ex);
 
-
-                    // todo 异常信息入库
+                    // 异常信息入库
+                    insertRecord("F", data);
                 } else {
                     log.info("同步数据发送到kafka消息系统成功,data: " + jsonMessage);
+                    insertRecord("S", data);
                 }
             });
         }catch (Exception e){
             log.error("同步数据发送到kafka消息系统异常,data: " + data, e);
-            // todo 异常信息入库
+            insertRecord("F", data);
         }
 
     }
 
+    //记录入库
+    public void insertRecord(String result, KafkaMessage<?> data){
+        try{
+            KafkaHeader header = data.getHeader();
+            String eventId = header.getEventId();
+            String sender = header.getSender();
+            String eventType = header.getEventType();
+            String tenantId = header.getTenantId();
+            // 信息入库
+            SendMessageRecordBo bo = new SendMessageRecordBo();
+            bo.setEventId(eventId);
+            bo.setSender(sender);
+            bo.setEventType(eventType);
+            bo.setTenantId(tenantId);
+            bo.setResult(result);
+            bo.setMessage(JSON.toJSONString(data));
+            LoginUser loginUser = LoginHelper.getLoginUser();
+            bo.setCreateBy(loginUser.getUserId());
+            bo.setMqType("kafka");
+            bo.setTopic(KafkaTopicConstants.SYNC_DATA_TOPIC);
+            bo.setTenantId(tenantId);
+            bo.setCreateTime(new Date());
+            sendMessageRecordService.insertByBo(bo);
+        }catch (Exception e){
+            log.error("kafka消息记录入库异常,data: " + data, e);
+        }
+    }
+
 }

+ 20 - 0
ruoyi-modules/ruoyi-backstage/src/main/resources/mapper/basics/SendMessageRecordMapper.xml

@@ -0,0 +1,20 @@
+<?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="org.dromara.backstage.basics.mapper.SendMessageRecordMapper">
+
+    <resultMap type="org.dromara.backstage.basics.domain.SendMessageRecord" id="SendMessageRecordResult">
+            <result property="recordId"    column="record_id"    />
+            <result property="mqType"    column="mq_type"    />
+            <result property="tenantId"    column="tenant_id"    />
+            <result property="topic"    column="topic"    />
+            <result property="eventType"    column="event_type"    />
+            <result property="result"    column="result"    />
+            <result property="createBy"    column="create_by"    />
+            <result property="createTime"    column="create_time"    />
+            <result property="message"    column="message"    />
+            <result property="eventId"    column="event_id"    />
+            <result property="sender"    column="sender"    />
+    </resultMap>
+</mapper>