You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SpringBoot接收PubSub消息存DB时主键冲突问题求助

问题描述

我正在实现一个简单场景:在应用中接收PubSub消息并将其保存到数据库。问题是可能会收到两条ID相同的消息,此时第二条消息应更新数据,但却遇到了主键约束冲突。

SQL执行日志如下:

Hibernate: 
    select
        batch0_.id as id1_0_0_,
        batch0_.date as date2_0_0_,
        batch0_.message as message3_0_0_,
        batch0_.status as status4_0_0_,
        batch0_.type as type5_0_0_ 
    from
        batch batch0_ 
    where
        batch0_.id=?
Hibernate: 
    select
        batch0_.id as id1_0_0_,
        batch0_.date as date2_0_0_,
        batch0_.message as message3_0_0_,
        batch0_.status as status4_0_0_,
        batch0_.type as type5_0_0_ 
    from
        batch batch0_ 
    where
        batch0_.id=?
Hibernate: 
    insert 
    into
        batch
        (date, message, status, type, id) 
    values
        (?, ?, ?, ?, ?)
Hibernate: 
    insert 
    into
        batch
        (date, message, status, type, id) 
    values
        (?, ?, ?, ?, ?)

看起来Hibernate在插入前对两条消息都执行了查询检查ID是否存在,但最终还是都执行了插入操作。

消息通过Spring服务激活器接收:

@ServiceActivator(inputChannel = "pubSubCheckConsistencyChannel")
public void messageReceiver(
      String payload,
      @Header(GcpPubSubHeaders.ORIGINAL_MESSAGE) BasicAcknowledgeablePubsubMessage message) {
    log.info("Message arrived on pubSubCheckConsistencyChannel Payload: " + payload);
    try {
      var batch = getBatchFromMessage(payload);
      saveBatch(batch);
      log.info("message  {} saved:", batch.getId());
    } catch (JsonProcessingException e) {
      log.error("error when try to read json message", e);
    } finally {
      message.ack();
    }
}

消息保存逻辑:

@Transactional
public void saveBatch(Batch batch) {
batchCacheManager.getCache("batch").put(batch.getId(), batchMapper.toDTO(batch));
  batchRepository.saveAndFlush(batch);
}

仓库定义:

@Repository
public interface BatchRepository extends JpaRepository<Batch, UUID>

实体类:

@Entity
@Getter
@Builder
@Setter
@NoArgsConstructor
@AllArgsConstructor
@ToString(callSuper = true, onlyExplicitlyIncluded = true)
@EqualsAndHashCode(onlyExplicitlyIncluded = true)
public class Batch  {
    @Id
    @Column(columnDefinition = "uuid",name = "id", nullable = false, unique = true)
    @ToString.Include
    @EqualsAndHashCode.Include
    private UUID id;
}

猜测是多线程导致的问题,想知道如何让服务激活器等待消息处理完成,或者如何同步插入操作。

解决方案

1. 限制PubSub消费者并发数

Spring Cloud GCP PubSub默认多线程消费消息,导致相同ID的消息被同时处理。通过配置将消费者线程数设为1,确保消息串行处理:

在application.properties中添加:

spring.cloud.gcp.pubsub.subscriber.max-concurrency=1

或application.yml中添加:

spring:
  cloud:
    gcp:
      pubsub:
        subscriber:
          max-concurrency: 1

2. 优化JPA保存逻辑

JpaRepository的save()方法依赖实体状态判断执行插入/更新,但并发场景下会出现"查询-插入"的竞态问题。改为先查询后更新的显式逻辑:

修改saveBatch方法:

@Transactional
public void saveBatch(Batch batch) {
    batchCacheManager.getCache("batch").put(batch.getId(), batchMapper.toDTO(batch));
    Optional<Batch> existingBatch = batchRepository.findById(batch.getId());
    if (existingBatch.isPresent()) {
        Batch updateBatch = existingBatch.get();
        updateBatch.setDate(batch.getDate());
        updateBatch.setMessage(batch.getMessage());
        updateBatch.setStatus(batch.getStatus());
        updateBatch.setType(batch.getType());
        batchRepository.saveAndFlush(updateBatch);
    } else {
        batchRepository.saveAndFlush(batch);
    }
}

3. 数据库层面加锁

悲观锁方案

在查询时锁定目标数据,避免并发修改:

首先修改Repository添加带锁的查询方法:

@Repository
public interface BatchRepository extends JpaRepository<Batch, UUID> {
    @Lock(LockModeType.PESSIMISTIC_WRITE)
    Optional<Batch> findById(UUID id);
}

再使用该方法执行保存逻辑(同方案2的saveBatch代码即可),悲观锁会在查询时锁定行数据,其他线程需等待锁释放后才能操作。

乐观锁方案

在实体类中添加版本字段,通过版本号控制更新冲突:

修改实体类:

@Entity
@Getter
@Builder
@Setter
@NoArgsConstructor
@AllArgsConstructor
@ToString(callSuper = true, onlyExplicitlyIncluded = true)
@EqualsAndHashCode(onlyExplicitlyIncluded = true)
public class Batch  {
    @Id
    @Column(columnDefinition = "uuid",name = "id", nullable = false, unique = true)
    @ToString.Include
    @EqualsAndHashCode.Include
    private UUID id;
    
    @Version
    private Integer version;
    // 其他原有字段
}

当并发更新时,版本号不匹配的操作会抛出OptimisticLockingFailureException,可捕获异常并重试。

4. 单实例同步块控制

在saveBatch方法中,以消息ID作为锁对象,确保同一ID的消息串行处理:

修改saveBatch方法:

@Transactional
public void saveBatch(Batch batch) {
    batchCacheManager.getCache("batch").put(batch.getId(), batchMapper.toDTO(batch));
    // 使用intern()保证相同ID的字符串引用同一对象,锁生效
    synchronized (batch.getId().toString().intern()) {
        Optional<Batch> existingBatch = batchRepository.findById(batch.getId());
        if (existingBatch.isPresent()) {
            Batch updateBatch = existingBatch.get();
            updateBatch.setDate(batch.getDate());
            updateBatch.setMessage(batch.getMessage());
            updateBatch.setStatus(batch.getStatus());
            updateBatch.setType(batch.getType());
            batchRepository.saveAndFlush(updateBatch);
        } else {
            batchRepository.saveAndFlush(batch);
        }
    }
}

注意:该方案仅适用于单实例应用,分布式场景下无效。

内容的提问来源于stack exchange,提问作者tschi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 18:00:00