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

使用@Transactional在插入新数据前覆盖PostgreSQL旧数据的事务问题

问题描述

加载包含相同日期的新数据集时,需要先删除PostgreSQL数据库中对应日期的所有现有记录。尝试使用带有REQUIRES_NEW/NESTED传播属性的@Transactional注解,但删除操作的事务并未在插入操作前提交,导致插入时触发唯一约束冲突。需要明确@Transactional的执行流程,找出逻辑问题,确保overwriteBpsData()的事务在saveAllBpsPositionData()的插入语句执行前完成提交。

代码示例
@Transactional
public void saveAllBpsPositionData(InputStream is) throws IOException {

    log.info("Parsing position data...");
    BPSParsingResult bpsParsingResult = bpsPositionParser.parse(is);
    if (!bpsParsingResult.getBpsPositions().isEmpty()) {
      LocalDate businessDate = bpsParsingResult.getBpsPositions().get(0).getBusinessDate();
      overwriteBpsData(businessDate);
    }
    log.info("Saving BPS Position Data...");
    if (!bpsParsingResult.getBpsPositions().isEmpty()) {
      insertBpsPositionData(bpsParsingResult.getBpsPositions());
    }
    log.info("Saving BPS Price Data...");
    if (!bpsParsingResult.getBpsPrices().isEmpty()) {
      insertBpsPriceData(bpsParsingResult.getBpsPrices());
    }
}

private void insertBpsPositionData(List<BPSPositionTable> bpsPositions) {
    try (Connection connection = hikariDataSource.getConnection();
        PreparedStatement positionStatement = connection.prepareStatement(POSITION_SQL);
        PreparedStatement memoStatement = connection.prepareStatement(MEMO_SQL)) {
      int count = 0;
      for (BPSPositionTable position : bpsPositions) {
        positionStatement.clearParameters();
        positionStatement.setString(1, position.getSRMASKey());
        // ....
        // ....
        positionStatement.setInt(33, position.getWhenIssueIndicator());
        positionStatement.setDate(
            34,
            position.getBusinessDate() == null ? null : Date.valueOf(position.getBusinessDate()));
        positionStatement.addBatch();

        if ((count + 1) % batchSize == 0 || (count + 1) == bpsPositions.size()) {
          positionStatement.executeBatch();
          positionStatement.clearBatch();
          memoStatement.executeBatch();
          memoStatement.clearBatch();
        }

        if (position.getNumberOfMemos() > 0) {
          for (BPSMemoTable memoTable : position.getCorrespondingMemos()) {
            memoStatement.setString(1, memoTable.getSRMASKey());
            memoStatement.setString(2, memoTable.getMemoTypeIndicator());
            memoStatement.setBigDecimal(3, memoTable.getMemoTypeQuantity());
            memoStatement.setDate(
                4,
                memoTable.getBusinessDate() == null
                    ? null
                    : Date.valueOf(memoTable.getBusinessDate()));
            memoStatement.addBatch();
          }
        }
        count++;
      }
      log.info("BPS Position Data Saved Successfully!");
    } catch (Exception e) {
      log.warn("Failure Inserting BPS Position Data: {}", e.getMessage());
    }
}

private void insertBpsPriceData(List<BPSPriceTable> bpsPrices) {
    try (Connection connection = hikariDataSource.getConnection();
        PreparedStatement statement = connection.prepareStatement(PRICE_SQL)) {
      int count = 0;
      for (BPSPriceTable price : bpsPrices) {
        statement.setInt(1, price.getClientNumber());
        statement.setString(2, price.getCusip());
        statement.setString(3, price.getSymbol());
        statement.setString(4, price.getCurrency());
        // ...
        // ...
        statement.setString(15, price.getSedol());
        statement.setDate(
            16, price.getBusinessDate() == null ? null : Date.valueOf(price.getBusinessDate()));
        statement.addBatch();
        if ((count + 1) % batchSize == 0 || (count + 1) == bpsPrices.size()) {
          statement.executeBatch();
          statement.clearBatch();
        }
        count++;
      }
      log.info("BPS Price Data Saved Successfully!");
    } catch (Exception e) {
      log.warn("Failure Inserting BPS Price Data: {}", e.getMessage());
    }
}

@Transactional(propagation = Propagation.REQUIRES_NEW) // Also have tried NESTED
public void overwriteBpsData(LocalDate businessDate) {
    if (bpsMemoRepo.countByBusinessDate(businessDate) > 0) {
      log.warn(
              "BPS Memo record(s) found by {} business date. Existing data will be overridden.",
              businessDate);
      bpsMemoRepo.deleteByBusinessDate(businessDate);
    }
    if (bpsPriceRepo.countByBusinessDate(businessDate) > 0) {
      log.warn(
              "BPS Price record(s) found by {} business date. Existing data will be overridden.",
              businessDate);
      bpsPriceRepo.deleteByBusinessDate(businessDate);
    }
    if (bpsPositionRepo.countByBusinessDate(businessDate) > 0) {
      log.warn(
              "BPS Position record(s) found by {} business date. Existing data will be overridden.",
              businessDate);
      bpsPositionRepo.deleteByBusinessDate(businessDate);
    }
}
问题分析与解决方案

核心问题1:内部方法调用导致事务注解失效

Spring的@Transactional基于动态代理实现,同一个类内部调用带注解的方法时,不会经过代理对象,因此overwriteBpsData()上的REQUIRES_NEW属性完全不生效。删除操作会直接加入saveAllBpsPositionData()的主事务,所有操作要等到主方法结束才提交,插入时旧数据未被删除,触发约束冲突。

核心问题2:原生JDBC连接脱离Spring事务管理

insertBpsPositionData()和insertBpsPriceData()直接从数据源获取连接,该连接不在Spring事务上下文内。即使主方法有事务,插入操作的事务和主事务相互独立——主事务的删除未提交,插入事务看不到删除状态,同样会引发冲突。

解决方案

方案1:拆分事务方法到独立Bean

把overwriteBpsData()抽取到单独的Spring管理Bean中,原服务注入该Bean后调用。此时调用经过代理对象,REQUIRES_NEW生效,删除操作在独立事务中执行并提交,之后再执行插入。

示例:

// 独立清理服务Bean
@Service
public class BpsDataCleanService {
    @Autowired
    private BpsMemoRepo bpsMemoRepo;
    @Autowired
    private BpsPriceRepo bpsPriceRepo;
    @Autowired
    private BpsPositionRepo bpsPositionRepo;

    @Transactional(propagation = Propagation.REQUIRES_NEW)
    public void overwriteBpsData(LocalDate businessDate) {
        // 原删除逻辑
    }
}

// 原服务中调用
@Service
public class BpsDataService {
    @Autowired
    private BpsDataCleanService cleanService;

    @Transactional
    public void saveAllBpsPositionData(InputStream is) throws IOException {
        // ...解析数据...
        cleanService.overwriteBpsData(businessDate);
        // ...执行插入...
    }
}

方案2:修改插入方法使用Spring事务连接

改用Spring的JdbcTemplate,或通过TransactionSynchronizationManager获取当前事务连接,确保插入操作加入Spring事务上下文:

示例(获取当前事务连接):

private void insertBpsPositionData(List<BPSPositionTable> bpsPositions) {
    Connection connection = DataSourceUtils.getConnection(hikariDataSource);
    try (PreparedStatement positionStatement = connection.prepareStatement(POSITION_SQL);
         PreparedStatement memoStatement = connection.prepareStatement(MEMO_SQL)) {
        // 原插入逻辑
    } catch (Exception e) {
        log.warn("Failure Inserting BPS Position Data: {}", e.getMessage());
        DataSourceUtils.releaseConnection(connection, hikariDataSource);
    }
}

配合方案1的独立删除事务,可保证删除提交后再执行插入。

方案3:(不推荐)使用AopContext调用代理方法

若不想拆分Bean,可通过AopContext.currentProxy()获取当前类代理对象调用方法:

@Transactional
public void saveAllBpsPositionData(InputStream is) throws IOException {
    // ...解析数据...
    ((YourServiceClass) AopContext.currentProxy()).overwriteBpsData(businessDate);
    // ...执行插入...
}

需在启动类添加@EnableAspectJAutoProxy(exposeProxy = true)开启代理暴露,该方式耦合Spring AOP实现,不推荐生产环境使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 12:15:48