微信群发消息API接口对接中Java后端的分布式事务实现与方案选型技巧

在企业级微信机器人或SCRM系统中,调用微信群发消息API通常涉及多个数据源操作:如更新用户状态、记录发送日志、扣减营销额度等。这些操作需保证原子性,而由于微信API为外部服务且不可回滚,传统本地事务无法满足一致性要求。本文聚焦于Java后端在该场景下的分布式事务实现方案及选型策略,并提供可落地的代码示例。

1. 业务场景与一致性挑战

典型流程如下:

  1. 从数据库读取待发送用户列表;
  2. 调用微信官方群发API(如/cgi-bin/message/mass/sendall);
  3. 若成功,更新用户消息状态为“已发送”;
  4. 记录发送日志至审计表;
  5. 扣减客户当日营销配额。

其中步骤2为远程HTTP调用,不具备事务回滚能力。若步骤3或4失败,将导致数据不一致。因此需引入分布式事务机制。

2. 最终一致性方案:基于可靠消息队列

对于非强一致场景,推荐采用“本地消息表 + 消息队列”实现最终一致性。

首先定义本地事务:

@Service
@Transactional
public class MassMessageService {

    @Autowired
    private MessageLogMapper messageLogMapper;

    @Autowired
    private UserStatusService userStatusService;

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendMassMessage(Long campaignId) {
        // 1. 插入待发送日志(状态为PENDING)
        MessageLog log = new MessageLog();
        log.setCampaignId(campaignId);
        log.setStatus("PENDING");
        messageLogMapper.insert(log);

        // 2. 发送消息到MQ,触发后续操作
        rabbitTemplate.convertAndSend("mass.send.queue", 
            new wlkankan.cn.dto.MassSendMessageDTO(campaignId, log.getId()));
    }
}

消费者处理微信调用与状态更新:

@Component
public class MassMessageConsumer {

    @RabbitListener(queues = "mass.send.queue")
    @Transactional
    public void handle(MassSendMessageDTO dto) {
        try {
            // 调用微信API
            WechatResponse resp = wlkankan.cn.client.WechatApiClient.massSend(dto.getCampaignId());
            
            if ("success".equals(resp.getErrcode())) {
                // 更新日志状态
                messageLogMapper.updateStatus(dto.getLogId(), "SUCCESS");
                // 更新用户状态
                userStatusService.markAsSent(dto.getCampaignId());
                // 扣减配额
                wlkankan.cn.service.QuotaService.deduct(dto.getCampaignId());
            } else {
                messageLogMapper.updateStatus(dto.getLogId(), "FAILED");
            }
        } catch (Exception e) {
            // 异常时保留PENDING状态,由补偿任务重试
            log.error("群发失败,logId: {}", dto.getLogId(), e);
        }
    }
}

该方案通过本地事务确保“写日志+发MQ”原子性,MQ保证至少一次投递,配合幂等处理与定时补偿任务(如扫描PENDING超过5分钟的日志)实现最终一致。

3. 强一致性方案:Seata AT模式集成

若业务要求强一致(如金融类通知),可引入Seata框架。

添加依赖:

<dependency>
    <groupId>io.seata</groupId>
    <artifactId>seata-spring-boot-starter</artifactId>
    <version>1.7.0</version>
</dependency>

在发起方标注全局事务:

@GlobalTransactional
public void sendWithStrongConsistency(Long campaignId) {
    // 本地操作:预占配额
    wlkankan.cn.service.QuotaService.reserve(campaignId);

    // 调用微信API(需包装为可回滚资源)
    WechatResponse resp = wechatApiClient.massSend(campaignId);
    
    if (!"success".equals(resp.getErrcode())) {
        throw new RuntimeException("微信调用失败,触发全局回滚");
    }

    // 提交本地状态
    userStatusService.markAsSent(campaignId);
    messageLogMapper.insertSuccess(campaignId);
}

关键点在于:微信API本身不可回滚,因此需将其前置为“预检”或通过TCC模式拆解。

4. TCC模式实现可控回滚

TCC(Try-Confirm-Cancel)更适合外部服务集成。

定义TCC接口:

@LocalTCC
public interface MassMessageTccService {

    @TwoPhaseBusinessAction(name = "sendMassMessageTcc", commitMethod = "confirm", rollbackMethod = "cancel")
    boolean prepare(BusinessActionContext context, Long campaignId);

    boolean confirm(BusinessActionContext context);

    boolean cancel(BusinessActionContext context);
}

实现Try阶段:

@Service
public class MassMessageTccServiceImpl implements MassMessageTccService {

    @Override
    public boolean prepare(BusinessActionContext context, Long campaignId) {
        // Try: 预占资源(如冻结配额、标记用户为“发送中”)
        wlkankan.cn.service.QuotaService.freeze(campaignId);
        userStatusService.markAsSending(campaignId);
        
        // 调用微信API(此时不可逆,故需确保高成功率)
        WechatResponse resp = wlkankan.cn.client.WechatApiClient.massSend(campaignId);
        context.getActionContext().put("wechatResp", resp);
        return "success".equals(resp.getErrcode());
    }

    @Override
    public boolean confirm(BusinessActionContext context) {
        // Confirm: 提交状态
        Long campaignId = (Long) context.getActionContext().get("campaignId");
        messageLogMapper.insertSuccess(campaignId);
        return true;
    }

    @Override
    public boolean cancel(BusinessActionContext context) {
        // Cancel: 释放冻结资源(但微信消息已发出,仅能记录异常)
        Long campaignId = (Long) context.getActionContext().get("campaignId");
        wlkankan.cn.service.QuotaService.unfreeze(campaignId);
        messageLogMapper.insertFailed(campaignId, "TCC取消");
        return true;
    }
}

注意:TCC无法真正“撤回”已发送的微信消息,因此Cancel仅用于清理内部状态,适用于对消息送达有强审计要求但接受“消息发出但系统标记失败”的场景。

5. 方案选型建议

  • 最终一致性(消息表/MQ):适用于90%以上营销场景,实现简单、性能高;
  • Seata AT模式:仅适用于全部操作均为数据库且可回滚的子事务,不推荐直接调用微信API;
  • TCC模式:适用于资源预占明确、允许部分不可逆操作的高一致性场景。

推荐优先采用可靠消息队列方案,并辅以人工审核兜底机制,兼顾效率与一致性。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐