支付系统核心流程详解

结合工作空间中使用的技术栈(Redis、Kafka、MyBatis、Spring Boot),以下是支付系统核心流程的详细实现方案及代码示例:


1. 请求接入与安全校验
接口鉴权实现
@RestController
@RequestMapping("/api/pay")
public class PaymentController {

    @Autowired
    private AuthService authService;
    
    @Autowired
    private PaymentService paymentService;

    @PostMapping("/create")
    public ResponseEntity<PaymentResult> createPayment(
            @RequestHeader("Authorization") String token,
            @Valid @RequestBody PaymentRequest request) {
        // 1. 接口鉴权
        boolean isAuthenticated = authService.validateToken(token);
        if (!isAuthenticated) {
            return ResponseEntity.status(HttpStatus.UNAUTHORIZED).body(
                new PaymentResult("FAIL", "身份验证失败"));
        }
        
        // 2. 风控检查
        RiskResult riskResult = authService.checkRisk(request);
        if (!riskResult.isPass()) {
            return ResponseEntity.badRequest().body(
                new PaymentResult("FAIL", riskResult.getReason()));
        }
        
        // 3. 处理支付请求
        PaymentResult result = paymentService.processPayment(request);
        return ResponseEntity.ok(result);
    }
}
风控服务实现
@Service
public class AuthService {
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;
    
    public boolean validateToken(String token) {
        // 从Redis中验证token
        String userId = (String) redisTemplate.opsForValue().get("token:" + token);
        return userId != null;
    }
    
    public RiskResult checkRisk(PaymentRequest request) {
        // 1. 检查高频交易
        String freqKey = "payment:freq:" + request.getUserId() + ":" + LocalDate.now();
        Long count = redisTemplate.opsForValue().increment(freqKey);
        if (count == 1) {
            redisTemplate.expire(freqKey, 1, TimeUnit.DAYS);
        }
        if (count > 10) { // 单日交易限额10笔
            return new RiskResult(false, "交易频率过高,请稍后再试");
        }
        
        // 2. 检查大额交易
        if (request.getAmount() > 10000) { // 单笔限额10000元
            return new RiskResult(false, "交易金额超出限制");
        }
        
        return new RiskResult(true, "风控检查通过");
    }
}

2. 核心业务流程
支付服务核心实现
@Service
public class PaymentService {
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;
    
    @Autowired
    private TransactionMapper transactionMapper;
    
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
    
    @Autowired
    private AccountService accountService;

    /**
     * 处理支付请求
     */
    @Transactional(rollbackFor = Exception.class)
    public PaymentResult processPayment(PaymentRequest request) {
        // 1. 参数校验
        if (!validateParams(request)) {
            return new PaymentResult("FAIL", "参数校验失败");
        }
        
        // 2. 幂等性检查
        String txnId = request.getTransactionId();
        if (checkIdempotency(txnId)) {
            return new PaymentResult("FAIL", "重复支付请求");
        }
        
        // 3. 账户余额校验
        boolean hasSufficientFunds = accountService.checkBalance(
            request.getUserId(), request.getAmount());
        if (!hasSufficientFunds) {
            return new PaymentResult("FAIL", "账户余额不足");
        }
        
        // 4. 分布式锁预扣款
        String lockKey = "lock:payment:" + request.getUserId();
        boolean locked = false;
        try {
            locked = redisTemplate.opsForValue().setIfAbsent(
                lockKey, "locked", 30, TimeUnit.SECONDS);
            
            if (!locked) {
                return new PaymentResult("FAIL", "系统繁忙,请稍后再试");
            }
            
            // 执行扣款
            boolean deductionSuccess = accountService.deductBalance(
                request.getUserId(), request.getAmount());
            if (!deductionSuccess) {
                throw new RuntimeException("扣款失败");
            }
            
            // 5. 生成交易流水
            TransactionRecord record = new TransactionRecord();
            record.setTransactionId(txnId);
            record.setUserId(request.getUserId());
            record.setAmount(request.getAmount());
            record.setStatus("SUCCESS");
            record.setCreateTime(new Date());
            transactionMapper.insert(record);
            
            // 6. 异步通知订单系统
            notifyOrderSystem(request, "SUCCESS");
            
            return new PaymentResult("SUCCESS", "支付成功");
        } catch (Exception e) {
            log.error("Payment processing failed: {}", e.getMessage(), e);
            // 7. 异常处理与补偿
            compensateDeduction(request.getUserId(), request.getAmount());
            notifyOrderSystem(request, "FAIL");
            return new PaymentResult("FAIL", "支付处理失败");
        } finally {
            // 释放锁
            if (locked) {
                redisTemplate.delete(lockKey);
            }
        }
    }
    
    /**
     * 参数校验
     */
    private boolean validateParams(PaymentRequest request) {
        return request != null
            && request.getUserId() != null
            && request.getAmount() > 0
            && StringUtils.hasText(request.getTransactionId())
            && StringUtils.hasText(request.getOrderId());
    }
    
    /**
     * 幂等性检查
     */
    private boolean checkIdempotency(String txnId) {
        String key = "payment:txn:" + txnId;
        Boolean exists = redisTemplate.hasKey(key);
        if (exists != null && exists) {
            return true;
        }
        // 设置过期时间,防止Redis内存溢出
        redisTemplate.opsForValue().set(key, "processed", 7, TimeUnit.DAYS);
        return false;
    }
    
    /**
     * 异步通知订单系统
     */
    private void notifyOrderSystem(PaymentRequest request, String status) {
        PaymentNotifyDTO notifyDTO = new PaymentNotifyDTO();
        notifyDTO.setOrderId(request.getOrderId());
        notifyDTO.setTransactionId(request.getTransactionId());
        notifyDTO.setStatus(status);
        notifyDTO.setAmount(request.getAmount());
        
        try {
            kafkaTemplate.send("payment-notify-topic", 
                JSON.toJSONString(notifyDTO));
            log.info("Sent payment notification: {}", notifyDTO);
        } catch (Exception e) {
            log.error("Failed to send payment notification: {}", e.getMessage(), e);
            // 可以将失败通知存入数据库,后续通过定时任务重试
        }
    }
    
    /**
     * 扣款补偿
     */
    private void compensateDeduction(Long userId, BigDecimal amount) {
        try {
            accountService.refundBalance(userId, amount);
            log.info("Compensated deduction for user: {}, amount: {}", userId, amount);
        } catch (Exception e) {
            log.error("Failed to compensate deduction: {}", e.getMessage(), e);
            // 记录补偿失败日志,后续人工处理
        }
    }
}
账户服务实现
@Service
public class AccountService {
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;
    
    @Autowired
    private AccountMapper accountMapper;

    /**
     * 检查账户余额
     */
    public boolean checkBalance(Long userId, BigDecimal amount) {
        // 优先查询Redis缓存
        String balanceKey = "account:balance:" + userId;
        BigDecimal balance = (BigDecimal) redisTemplate.opsForValue().get(balanceKey);
        
        if (balance == null) {
            // 缓存未命中,查询数据库
            Account account = accountMapper.selectById(userId);
            if (account == null) {
                return false;
            }
            balance = account.getBalance();
            // 更新缓存
            redisTemplate.opsForValue().set(balanceKey, balance);
        }
        
        return balance.compareTo(amount) >= 0;
    }
    
    /**
     * 扣减账户余额
     */
    @Transactional(rollbackFor = Exception.class)
    public boolean deductBalance(Long userId, BigDecimal amount) {
        // 更新数据库
        int updated = accountMapper.deductBalance(userId, amount);
        if (updated <= 0) {
            return false;
        }
        
        // 更新缓存(延迟双删策略)
        String balanceKey = "account:balance:" + userId;
        redisTemplate.delete(balanceKey);
        
        // 异步加载最新余额到缓存
        CompletableFuture.runAsync(() -> {
            try {
                Thread.sleep(100); // 等待数据库事务提交
                Account account = accountMapper.selectById(userId);
                redisTemplate.opsForValue().set(balanceKey, account.getBalance());
            } catch (Exception e) {
                log.error("Failed to refresh balance cache: {}", e.getMessage(), e);
            }
        });
        
        return true;
    }
    
    /**
     * 退还账户余额
     */
    @Transactional(rollbackFor = Exception.class)
    public boolean refundBalance(Long userId, BigDecimal amount) {
        accountMapper.refundBalance(userId, amount);
        
        // 失效缓存
        String balanceKey = "account:balance:" + userId;
        redisTemplate.delete(balanceKey);
        return true;
    }
}

3. 数据持久化
MyBatis映射文件
<!-- TransactionMapper.xml -->
<mapper namespace="com.bigmarket.mapper.TransactionMapper">
    <insert id="insert" parameterType="com.bigmarket.entity.TransactionRecord">
        INSERT INTO transaction_record (
            transaction_id, user_id, amount, status, create_time
        ) VALUES (
            #{transactionId}, #{userId}, #{amount}, #{status}, #{createTime}
        )
    </insert>
</mapper>

<!-- AccountMapper.xml -->
<mapper namespace="com.bigmarket.mapper.AccountMapper">
    <select id="selectById" parameterType="java.lang.Long" resultType="com.bigmarket.entity.Account">
        SELECT id, user_id, balance FROM account WHERE user_id = #{id}
    </select>
    
    <update id="deductBalance" parameterType="map">
        UPDATE account
        SET balance = balance - #{amount}
        WHERE user_id = #{userId} AND balance >= #{amount}
    </update>
    
    <update id="refundBalance" parameterType="map">
        UPDATE account
        SET balance = balance + #{amount}
        WHERE user_id = #{userId}
    </update>
</mapper>

4. 异常处理与监控
全局异常处理
@RestControllerAdvice
public class GlobalExceptionHandler {
    @ExceptionHandler(Exception.class)
    public ResponseEntity<ErrorResponse> handleException(Exception e) {
        log.error("Global exception caught: {}", e.getMessage(), e);
        ErrorResponse error = new ErrorResponse(
            "INTERNAL_ERROR", "系统内部错误,请稍后再试");
        return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body(error);
    }
}
监控指标收集
@Service
public class MetricsService {
    @Autowired
    private MeterRegistry meterRegistry;
    
    private Counter paymentSuccessCounter;
    private Counter paymentFailureCounter;
    private Timer paymentTimer;
    
    @PostConstruct
    public void init() {
        paymentSuccessCounter = Counter.builder("payment.success.count")
            .description("Number of successful payments")
            .register(meterRegistry);
        
        paymentFailureCounter = Counter.builder("payment.failure.count")
            .description("Number of failed payments")
            .register(meterRegistry);
        
        paymentTimer = Timer.builder("payment.processing.time")
            .description("Payment processing time")
            .register(meterRegistry);
    }
    
    public void recordSuccess() {
        paymentSuccessCounter.increment();
    }
    
    public void recordFailure() {
        paymentFailureCounter.increment();
    }
    
    public <T> T recordProcessingTime(Supplier<T> supplier) {
        return paymentTimer.record(supplier);
    }
}

5. 关键配置
Redis配置
@Configuration
@EnableCaching
public class RedisConfig {
    @Bean
    public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
        RedisTemplate<String, Object> template = new RedisTemplate<>();
        template.setConnectionFactory(factory);
        
        // 使用Jackson2JsonRedisSerializer序列化值
        Jackson2JsonRedisSerializer<Object> serializer = new Jackson2JsonRedisSerializer<>(Object.class);
        ObjectMapper mapper = new ObjectMapper();
        mapper.setVisibility(PropertyAccessor.ALL, JsonAutoDetect.Visibility.ANY);
        mapper.activateDefaultTyping(LaissezFaireSubTypeValidator.instance, 
            ObjectMapper.DefaultTyping.NON_FINAL);
        serializer.setObjectMapper(mapper);
        
        template.setValueSerializer(serializer);
        template.setKeySerializer(new StringRedisSerializer());
        template.setHashKeySerializer(new StringRedisSerializer());
        template.setHashValueSerializer(serializer);
        template.afterPropertiesSet();
        
        return template;
    }
}
Kafka配置
@Configuration
public class KafkaConfig {
    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;
    
    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        // 开启事务支持
        configProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "payment-transaction");
        return new DefaultKafkaProducerFactory<>(configProps);
    }
    
    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        KafkaTemplate<String, String> template = new KafkaTemplate<>(producerFactory());
        template.setTransactionIdPrefix("payment-");
        return template;
    }
}

总结

以上代码示例完整展示了支付系统的核心流程,包括:

  1. 安全接入:接口鉴权和风控检查
  2. 核心业务:参数校验、幂等性检查、账户余额校验、分布式锁预扣款、交易流水生成、异步通知
  3. 异常处理:事务回滚、扣款补偿、失败通知
  4. 数据一致性:Redis缓存一致性策略、分布式事务支持
  5. 监控告警:指标收集、日志记录

该实现充分利用了工作空间中提到的技术栈(Redis、Kafka、MyBatis、Spring Boot),并解决了支付系统中的关键问题(如分布式锁、幂等性、数据一致性等)。实际项目中,还可以根据业务需求进一步扩展,如添加多支付渠道支持、更复杂的风控规则、分库分表等。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐