Browse Source

feat(网体): 新增网体统计功能

lw 2 months ago
parent
commit
caac9a3626

+ 23 - 0
bex-cloud-staking-core/src/main/java/com/bex/staking/model/dto/NetworkPeriodAmountDTO.java

@@ -0,0 +1,23 @@
+package com.bex.staking.model.dto;
+
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.math.BigDecimal;
+
+/**
+ * 网体统计 - 按周期分桶聚合金额结果。
+ *
+ * <p>periodKey 含义随 period 参数变化:
+ * DAY -> yyyy-MM-dd;WEEK -> 该自然周周一的日期 yyyy-MM-dd;MONTH -> yyyy-MM。</p>
+ */
+@Data
+@NoArgsConstructor
+public class NetworkPeriodAmountDTO {
+
+    /** 周期 key */
+    private String periodKey;
+
+    /** 该周期累计金额 */
+    private BigDecimal amount;
+}

+ 23 - 0
bex-cloud-staking-core/src/main/java/com/bex/staking/model/dto/NetworkUserAmountDTO.java

@@ -0,0 +1,23 @@
+package com.bex.staking.model.dto;
+
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.math.BigDecimal;
+
+/**
+ * 网体统计 - 按用户聚合金额结果。
+ *
+ * <p>仅供 Mapper 聚合查询使用,userId -> amount 的扁平结构,
+ * 由 Service 层转换为 {@code Map<Long, BigDecimal>}。</p>
+ */
+@Data
+@NoArgsConstructor
+public class NetworkUserAmountDTO {
+
+    /** 用户 id */
+    private Long userId;
+
+    /** 累计金额 */
+    private BigDecimal amount;
+}

+ 57 - 0
bex-cloud-staking-entity/src/main/java/com/bex/staking/mapper/StakingOrderMapper.java

@@ -3,6 +3,9 @@ package com.bex.staking.mapper;
 import com.baomidou.mybatisplus.core.mapper.BaseMapper;
 import com.baomidou.mybatisplus.core.mapper.BaseMapper;
 import com.baomidou.mybatisplus.core.metadata.IPage;
 import com.baomidou.mybatisplus.core.metadata.IPage;
 import com.bex.staking.entity.StakingOrder;
 import com.bex.staking.entity.StakingOrder;
+import com.bex.staking.model.dto.NetworkPeriodAmountDTO;
+import com.bex.staking.model.dto.NetworkUserAmountDTO;
+import java.math.BigDecimal;
 import java.util.List;
 import java.util.List;
 import org.apache.ibatis.annotations.Mapper;
 import org.apache.ibatis.annotations.Mapper;
 import org.apache.ibatis.annotations.Param;
 import org.apache.ibatis.annotations.Param;
@@ -77,4 +80,58 @@ public interface StakingOrderMapper extends BaseMapper<StakingOrder> {
     List<StakingOrder> scanHoldingOrders(
     List<StakingOrder> scanHoldingOrders(
             @Param("lastId") Long lastId,
             @Param("lastId") Long lastId,
             @Param("limit") int limit);
             @Param("limit") int limit);
+
+    /**
+     * 网体统计 - 查询给定 userIds 中在 staking_order 表有记录(已参与质押)的用户 id。
+     *
+     * <p>仅用于"已参与人数"判定,结果去重;不参与逻辑删除过滤策略时仍显式追加 deleted = 0。</p>
+     *
+     * @param userIds 候选用户 id 列表(非空,调用方需分片,单批建议 ≤ 1000)
+     * @return 已参与用户 id 列表
+     */
+    List<Long> findParticipatedUserIds(@Param("userIds") List<Long> userIds);
+
+    /**
+     * 网体统计 - 按 userIds 累计 stake_amount 总和。
+     *
+     * <p>当 startTime / endTime 任一为 null 时,对应方向不做时间过滤(用于"累积交易量");
+     * 当两个时间都传入时,按 created_at 区间过滤(用于柱状图区间累计)。</p>
+     *
+     * @param userIds 候选用户 id 列表(非空,调用方需分片,单批建议 ≤ 1000)
+     * @param startTime 起始时间(含),可为 null 表示不限
+     * @param endTime 结束时间(含),可为 null 表示不限
+     * @return 总金额;无记录返回 null(调用方按 0 处理)
+     */
+    BigDecimal sumStakeAmountByUserIds(
+            @Param("userIds") List<Long> userIds,
+            @Param("startTime") java.time.LocalDateTime startTime,
+            @Param("endTime") java.time.LocalDateTime endTime);
+
+    /**
+     * 网体统计 - 按 userIds 分用户累计 stake_amount(接口3 的当页交易金额聚合)。
+     *
+     * @param userIds 候选用户 id 列表(非空,调用方需分片,单批建议 ≤ 1000)
+     * @return 用户 id -> 累计金额(无记录用户不返回)
+     */
+    List<NetworkUserAmountDTO> sumStakeAmountGroupByUser(@Param("userIds") List<Long> userIds);
+
+    /**
+     * 网体统计 - 按周期分桶聚合 stake_amount,用于柱状图。
+     *
+     * <p>period 取值:DAY / WEEK / MONTH。
+     * - DAY:DATE_FORMAT(created_at,'%Y-%m-%d');
+     * - WEEK:按周一为首日,DATE_FORMAT(DATE_SUB(created_at, INTERVAL WEEKDAY(created_at) DAY),'%Y-%m-%d');
+     * - MONTH:DATE_FORMAT(created_at,'%Y-%m')。</p>
+     *
+     * @param userIds 候选用户 id 列表(非空,调用方需分片,单批建议 ≤ 1000)
+     * @param period 周期类型 DAY/WEEK/MONTH
+     * @param startTime 起始时间(含)
+     * @param endTime 结束时间(含)
+     * @return 周期 key -> 金额
+     */
+    List<NetworkPeriodAmountDTO> sumStakeAmountByPeriod(
+            @Param("userIds") List<Long> userIds,
+            @Param("period") String period,
+            @Param("startTime") java.time.LocalDateTime startTime,
+            @Param("endTime") java.time.LocalDateTime endTime);
 }
 }

+ 71 - 0
bex-cloud-staking-entity/src/main/resources/mapper/StakingOrderMapper.xml

@@ -100,4 +100,75 @@
         LIMIT #{limit}
         LIMIT #{limit}
     </select>
     </select>
 
 
+    <!-- 网体统计 - 查询给定 userIds 中已参与质押的用户 id(去重) -->
+    <select id="findParticipatedUserIds" resultType="java.lang.Long">
+        SELECT DISTINCT user_id
+        FROM staking_order
+        WHERE deleted = 0
+          AND user_id IN
+          <foreach collection="userIds" item="uid" open="(" separator="," close=")">
+              #{uid}
+          </foreach>
+    </select>
+
+    <!-- 网体统计 - 按 userIds 累计 stake_amount,可选 created_at 时间区间过滤 -->
+    <select id="sumStakeAmountByUserIds" resultType="java.math.BigDecimal">
+        SELECT COALESCE(SUM(stake_amount), 0)
+        FROM staking_order
+        WHERE deleted = 0
+          AND user_id IN
+          <foreach collection="userIds" item="uid" open="(" separator="," close=")">
+              #{uid}
+          </foreach>
+        <if test="startTime != null">
+            AND created_at &gt;= #{startTime}
+        </if>
+        <if test="endTime != null">
+            AND created_at &lt;= #{endTime}
+        </if>
+    </select>
+
+    <!-- 网体统计 - 按 userIds 分用户累计 stake_amount -->
+    <select id="sumStakeAmountGroupByUser"
+            resultType="com.bex.staking.model.dto.NetworkUserAmountDTO">
+        SELECT user_id AS userId, COALESCE(SUM(stake_amount), 0) AS amount
+        FROM staking_order
+        WHERE deleted = 0
+          AND user_id IN
+          <foreach collection="userIds" item="uid" open="(" separator="," close=")">
+              #{uid}
+          </foreach>
+        GROUP BY user_id
+    </select>
+
+    <!--
+      网体统计 - 按周期分桶聚合 stake_amount,用于柱状图。
+      WEEK 以周一为首日:WEEKDAY() 周一返回 0,周日返回 6。
+    -->
+    <select id="sumStakeAmountByPeriod"
+            resultType="com.bex.staking.model.dto.NetworkPeriodAmountDTO">
+        SELECT
+        <choose>
+            <when test="period == 'DAY'">
+                DATE_FORMAT(created_at, '%Y-%m-%d') AS periodKey,
+            </when>
+            <when test="period == 'WEEK'">
+                DATE_FORMAT(DATE_SUB(created_at, INTERVAL WEEKDAY(created_at) DAY), '%Y-%m-%d') AS periodKey,
+            </when>
+            <otherwise>
+                DATE_FORMAT(created_at, '%Y-%m') AS periodKey,
+            </otherwise>
+        </choose>
+        COALESCE(SUM(stake_amount), 0) AS amount
+        FROM staking_order
+        WHERE deleted = 0
+          AND user_id IN
+          <foreach collection="userIds" item="uid" open="(" separator="," close=")">
+              #{uid}
+          </foreach>
+          AND created_at &gt;= #{startTime}
+          AND created_at &lt;= #{endTime}
+        GROUP BY periodKey
+    </select>
+
 </mapper>
 </mapper>

+ 143 - 0
bex-cloud-staking-service/src/main/java/com/bex/staking/grpc/StakingGrpcService.java

@@ -6,10 +6,16 @@ import com.bex.proto.staking.*;
 import com.bex.proto.user.*;
 import com.bex.proto.user.*;
 import com.bex.staking.entity.StakingOrder;
 import com.bex.staking.entity.StakingOrder;
 import com.bex.staking.entity.StakingProduct;
 import com.bex.staking.entity.StakingProduct;
+import com.bex.staking.entity.StakingProfitJournal;
 import com.bex.staking.enums.ProductStatus;
 import com.bex.staking.enums.ProductStatus;
 import com.bex.staking.mapper.StakingOrderMapper;
 import com.bex.staking.mapper.StakingOrderMapper;
 import com.bex.staking.mapper.StakingProductMapper;
 import com.bex.staking.mapper.StakingProductMapper;
+import com.bex.staking.mapper.StakingProfitJournalMapper;
+import com.bex.staking.model.dto.NetworkPeriodAmountDTO;
+import com.bex.staking.model.dto.NetworkUserAmountDTO;
 import com.bex.staking.util.ProductDisplayMapper;
 import com.bex.staking.util.ProductDisplayMapper;
+
+import java.time.Instant;
 import io.grpc.stub.StreamObserver;
 import io.grpc.stub.StreamObserver;
 import lombok.RequiredArgsConstructor;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import lombok.extern.slf4j.Slf4j;
@@ -35,6 +41,7 @@ public class StakingGrpcService extends StakingServiceGrpc.StakingServiceImplBas
     private final StakingProductMapper stakingProductMapper;
     private final StakingProductMapper stakingProductMapper;
     private final StakingOrderMapper stakingOrderMapper;
     private final StakingOrderMapper stakingOrderMapper;
     private final UserServiceGrpc.UserServiceBlockingStub userStub;
     private final UserServiceGrpc.UserServiceBlockingStub userStub;
+    private final StakingProfitJournalMapper stakingProfitJournalMapper;
 
 
     @Override
     @Override
     public void listProducts(ListProductsRequest request, StreamObserver<ListProductsResponse> responseObserver) {
     public void listProducts(ListProductsRequest request, StreamObserver<ListProductsResponse> responseObserver) {
@@ -377,6 +384,42 @@ public class StakingGrpcService extends StakingServiceGrpc.StakingServiceImplBas
         }
         }
     }
     }
 
 
+    /**
+     * 按 journal id 查询单笔收益派发流水(佣金对账专用)
+     * <p>
+     * 用于 commission 服务在 T+1 对账阶段重算质押返佣,输入参数为
+     * {@code commission_records.sourceTradeId}(即 {@code staking_profit_journal.id}),
+     * 返回该笔流水的 {@code profit_usdt_amount} 作为重算 sourceFee。
+     * </p>
+     */
+    @Override
+    public void getProfitJournal(GetProfitJournalRequest request, StreamObserver<GetProfitJournalResponse> responseObserver) {
+        try {
+            long journalId = request.getJournalId();
+            log.info("【质押对账】gRPC GetProfitJournal 收到请求:journalId={}", journalId);
+
+            StakingProfitJournal journal = stakingProfitJournalMapper.selectById(journalId);
+            GetProfitJournalResponse.Builder builder = GetProfitJournalResponse.newBuilder();
+            if (journal == null) {
+                log.info("【质押对账】gRPC GetProfitJournal 未找到流水:journalId={}", journalId);
+                builder.setFound(false).setProfitAmount("0").setUserId(0L);
+            } else {
+                BigDecimal profitUsdt = journal.getProfitAmount() != null
+                        ? journal.getProfitAmount() : BigDecimal.ZERO;
+                builder.setFound(true)
+                        .setProfitAmount(profitUsdt.toPlainString())
+                        .setUserId(journal.getUserId() != null ? journal.getUserId() : 0L);
+                log.info("【质押对账】gRPC GetProfitJournal 查询成功:journalId={}, userId={}, profitAmount={}",
+                        journalId, journal.getUserId(), profitUsdt);
+            }
+            responseObserver.onNext(builder.build());
+            responseObserver.onCompleted();
+        } catch (Exception e) {
+            log.error("【质押对账】gRPC GetProfitJournal 异常:journalId={}", request.getJournalId(), e);
+            responseObserver.onError(io.grpc.Status.INTERNAL.withDescription(e.getMessage()).asRuntimeException());
+        }
+    }
+
     private StakingProductInfo toProductInfo(StakingProduct p) {
     private StakingProductInfo toProductInfo(StakingProduct p) {
         StakingProductInfo.Builder builder = StakingProductInfo.newBuilder()
         StakingProductInfo.Builder builder = StakingProductInfo.newBuilder()
                 .setId(p.getId() != null ? p.getId() : 0L)
                 .setId(p.getId() != null ? p.getId() : 0L)
@@ -460,4 +503,104 @@ public class StakingGrpcService extends StakingServiceGrpc.StakingServiceImplBas
 
 
         return builder.build();
         return builder.build();
     }
     }
+
+    // ==================== 网体统计 RPC ====================
+
+    @Override
+    public void getStakingParticipatedUserIds(NetworkParticipatedUserIdsRequest request,
+                                              StreamObserver<NetworkParticipatedUserIdsResponse> responseObserver) {
+        try {
+            List<Long> userIds = request.getUserIdsList();
+            NetworkParticipatedUserIdsResponse.Builder builder = NetworkParticipatedUserIdsResponse.newBuilder();
+            if (userIds == null || userIds.isEmpty()) {
+                responseObserver.onNext(builder.build());
+                responseObserver.onCompleted();
+                return;
+            }
+            List<Long> participated = stakingOrderMapper.findParticipatedUserIds(userIds);
+            if (participated != null) {
+                builder.addAllParticipatedUserIds(participated);
+            }
+            responseObserver.onNext(builder.build());
+            responseObserver.onCompleted();
+        } catch (Exception e) {
+            log.error("getStakingParticipatedUserIds gRPC error: {}", e.getMessage(), e);
+            responseObserver.onError(io.grpc.Status.INTERNAL.withDescription(e.getMessage()).asRuntimeException());
+        }
+    }
+
+    @Override
+    public void sumStakingAmountByUserIds(NetworkSumAmountRequest request,
+                                          StreamObserver<NetworkSumAmountResponse> responseObserver) {
+        try {
+            List<Long> userIds = request.getUserIdsList();
+            NetworkSumAmountResponse.Builder builder = NetworkSumAmountResponse.newBuilder()
+                    .setTotalAmount("0");
+            if (userIds == null || userIds.isEmpty()) {
+                responseObserver.onNext(builder.build());
+                responseObserver.onCompleted();
+                return;
+            }
+            if (request.getGroupByUser()) {
+                List<NetworkUserAmountDTO> list = stakingOrderMapper.sumStakeAmountGroupByUser(userIds);
+                BigDecimal total = BigDecimal.ZERO;
+                if (list != null) {
+                    for (NetworkUserAmountDTO dto : list) {
+                        if (dto.getUserId() == null) continue;
+                        BigDecimal amt = dto.getAmount() != null ? dto.getAmount() : BigDecimal.ZERO;
+                        builder.putAmountByUser(dto.getUserId(), amt.toPlainString());
+                        total = total.add(amt);
+                    }
+                }
+                builder.setTotalAmount(total.toPlainString());
+            } else {
+                LocalDateTime startTime = request.getStartTimeMs() > 0
+                        ? LocalDateTime.ofInstant(Instant.ofEpochMilli(request.getStartTimeMs()), ZoneId.systemDefault())
+                        : null;
+                LocalDateTime endTime = request.getEndTimeMs() > 0
+                        ? LocalDateTime.ofInstant(Instant.ofEpochMilli(request.getEndTimeMs()), ZoneId.systemDefault())
+                        : null;
+                BigDecimal sum = stakingOrderMapper.sumStakeAmountByUserIds(userIds, startTime, endTime);
+                builder.setTotalAmount(sum != null ? sum.toPlainString() : "0");
+            }
+            responseObserver.onNext(builder.build());
+            responseObserver.onCompleted();
+        } catch (Exception e) {
+            log.error("sumStakingAmountByUserIds gRPC error: {}", e.getMessage(), e);
+            responseObserver.onError(io.grpc.Status.INTERNAL.withDescription(e.getMessage()).asRuntimeException());
+        }
+    }
+
+    @Override
+    public void sumStakingAmountByPeriod(NetworkSumAmountByPeriodRequest request,
+                                         StreamObserver<NetworkSumAmountByPeriodResponse> responseObserver) {
+        try {
+            List<Long> userIds = request.getUserIdsList();
+            NetworkSumAmountByPeriodResponse.Builder builder = NetworkSumAmountByPeriodResponse.newBuilder();
+            if (userIds == null || userIds.isEmpty()) {
+                responseObserver.onNext(builder.build());
+                responseObserver.onCompleted();
+                return;
+            }
+            LocalDateTime startTime = LocalDateTime.ofInstant(
+                    Instant.ofEpochMilli(request.getStartTimeMs()), ZoneId.systemDefault());
+            LocalDateTime endTime = LocalDateTime.ofInstant(
+                    Instant.ofEpochMilli(request.getEndTimeMs()), ZoneId.systemDefault());
+            List<NetworkPeriodAmountDTO> list = stakingOrderMapper.sumStakeAmountByPeriod(
+                    userIds, request.getPeriod(), startTime, endTime);
+            if (list != null) {
+                for (NetworkPeriodAmountDTO dto : list) {
+                    builder.addItems(NetworkPeriodAmount.newBuilder()
+                            .setPeriodKey(dto.getPeriodKey() != null ? dto.getPeriodKey() : "")
+                            .setAmount(dto.getAmount() != null ? dto.getAmount().toPlainString() : "0")
+                            .build());
+                }
+            }
+            responseObserver.onNext(builder.build());
+            responseObserver.onCompleted();
+        } catch (Exception e) {
+            log.error("sumStakingAmountByPeriod gRPC error: {}", e.getMessage(), e);
+            responseObserver.onError(io.grpc.Status.INTERNAL.withDescription(e.getMessage()).asRuntimeException());
+        }
+    }
 }
 }