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

