支付系统核心流程详解
·
支付系统核心流程详解
结合工作空间中使用的技术栈(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;
}
}
总结
以上代码示例完整展示了支付系统的核心流程,包括:
- 安全接入:接口鉴权和风控检查
- 核心业务:参数校验、幂等性检查、账户余额校验、分布式锁预扣款、交易流水生成、异步通知
- 异常处理:事务回滚、扣款补偿、失败通知
- 数据一致性:Redis缓存一致性策略、分布式事务支持
- 监控告警:指标收集、日志记录
该实现充分利用了工作空间中提到的技术栈(Redis、Kafka、MyBatis、Spring Boot),并解决了支付系统中的关键问题(如分布式锁、幂等性、数据一致性等)。实际项目中,还可以根据业务需求进一步扩展,如添加多支付渠道支持、更复杂的风控规则、分库分表等。
更多推荐
所有评论(0)