当前位置: 首页 > 知识库问答 >
问题:

r2dbc-pool连接取消后未释放

齐昊苍
2023-03-14

我对R2DBC池有一个奇怪的行为:我们碰巧创建了大量线程并将它们发送到R2DBC池以获取DB连接。当池中的所有R2DBC连接都在使用时,我们创建的线程排队等待空闲连接可用,这发生在释放之前使用的连接时。如果我们在等待空闲连接时取消这些线程,会发生以下行为:

  • 即使它们被取消,也有一些线程获得连接并通过其正常的DB进程

因此,一些连接不会返回到空闲状态。它们保持被获取状态,并阻止后续连接请求获取这些特定连接。重新启动服务之前,连接将保持锁定状态。

值得一提的是,我们在获取连接时对数据库进行了查询(我们有一个多租户数据库,并在获取连接时使用SET SCHEMA来选择正确的租户)。

我制作了一个程序来复制这个问题。

对于测试,我使用maxConnection=2的池。在多次调用测试方法(controller.test)之后,池中的一些连接仍然会被无限期地获取(它们应该都是由onCancel或由Spring处理的close语句释放的)。通过使用jmx监控池,可以很容易地证明这一点。

我假设取消请求传播到connectionPool。create(),但有些迭代似乎有足够的时间在收到取消之前结束预查询,这导致连接可供Spring使用。在这种情况下,在TestConnectionFactory中看不到取消,大约有1/3次,Spring不调用连接。关闭,导致连接仍被获取。

@Slf4j
@RestController
public class TestController {
    private final TestRepo1 testRepo1;

    @Autowired
    public TestController(
            TestRepo1 testRepo1
    ) {
        this.testRepo1 = testRepo1;
    }

    @GetMapping("test")
    Mono<Void> test(
    ) {
        // Will made 49 queries to the database.
        return Mono
                .when(
                        IntStream.range(0, 100)
                                .mapToObj(i -> Mono.defer(() ->
                                        i == 0 ? // the first element throw an error after 2 seconds, canceling all query not already done.
                                                Mono.just(0)
                                                        .delayElement(Duration.ofMillis(2000))
                                                        .doOnNext(x -> log.info("{} -> throw", x))
                                                        .then(Mono.error(new Exception("FAIL"))) :
                                                testRepo1.query(String.valueOf(i)))
                                )
                                .collect(Collectors.toList())
                )
                .then()
                .onErrorResume(e -> Mono.empty()); // avoid propagating error to http response.
    }
}
@Slf4j
public class TestConnectionFactory implements ConnectionFactory {
    private final ConnectionPool connectionPool;

    TestConnectionFactory(ConnectionPool connectionPool) {
        this.connectionPool = connectionPool;
    }

    @Override
    public Publisher<? extends Connection> create() {
       return createTenantConnection()
                .doOnNext(x -> log.info("creation transaction done"))
                .doOnCancel(() -> log.info("cancel while creation"));
    }

    private Mono<Connection> createTenantConnection() {
        return connectionPool.create()
                .flatMap(connection -> preQuery(connection));
    }

    private Mono<Connection> preQuery(Connection connection) {
        return Mono.from(connection
                .createStatement("SELECT 1;") // enough to produce the error, in our real code, this is a SET SCHEMA XXX
                .execute())
                .doOnCancel(() -> log.info("cancel during preQuery"))
                .thenReturn(connection);
    }

    @Override
    public ConnectionFactoryMetadata getMetadata() {
        return connectionPool.getMetadata();
    }
}
@Configuration
public class MyConfiguration {
    @Bean
    @Scope("singleton")
    ConnectionFactory connectionFactory(
            ConnectionPool connectionPool
    ) {
        return new TestConnectionFactory(connectionPool);
    }
}
@Slf4j
@Repository
public class TestRepo1 {
    // simple query waiting 1 second
    private static final String QUERY = "SELECT pg_sleep(1);";

    private final DatabaseClient databaseClient;

    @Autowired
    public TestRepo1(DatabaseClient databaseClient) {
        this.databaseClient = databaseClient;
    }

    public Mono<Void> query(String msg) {
        log.info("start query {}", msg);
        return databaseClient.execute(QUERY)
                .map(row -> "result")
                .first()
                .doOnCancel(() -> log.info("cancel query {}", msg))
                .doOnNext(x -> log.info("query {} result", msg))
                .then()
                .doOnTerminate(() -> log.info("terminate {}", msg));
    }
}

我们将org.springframework.boot2.3.5. RELEASE与io. r2dbc: r2dbc-postgresql和io. r2dbc: r2dbc-pool一起使用。

我们尝试升级到io. r2dbc: r2dbc-postgresql 0.8.8。RELEASE和io. r2dbc: r2dbc-pool 0.9.0。M1但结果保持不变。

共有2个答案

凤凡
2023-03-14

如果您使用Spring来管理r2dbc连接,您还应该使用Spring事务管理器来正确处理所有相关资源(连接、池等),至少在理论上是这样(实际上在生产中不太常见——几个月不到一次——我仍然会看到资源泄漏的情况——我认为这是因为问题:https://github.com/r2dbc/r2dbc-pool/issues/140——我应该更新r2dbc-pool版本)。同样的问题可能是您问题的原因。尽管如此,您应该正确关闭资源。否则连接会泄漏。我只知道这里有两个选项:

  1. 使用Flux.using当-如已建议的
  2. 使用Spring Transaction Manager

要使用Spring Transaction Manager,您应该创建它。最简单和最方便的方法是使用相应的Spring引导自动配置:请参阅org.springframework.boot中的类R2dbcTransactionManagerAutoConfiguration: spring-boot-autoconfiure库(您也可以使用R2dbcAutoConfiguration根据配置yaml文件而不是手动bean创建自动创建连接工厂)。

创建事务管理器后,您应该创建将封装存储库并将所需逻辑包装到事务上下文的服务。例如:

class Service {
  private final Repository repo;
  private final Repository2 repo2;
   ...
  @Transactional // R2DBC
  public Mono<Result> process(…) {
    return repo.save(e).zipWith(
      repo2.save(e2), (t, t2) -> t2);  
  }
}

如果需要,可以使用注释整个类或特定方法。

在您的示例中,您可以只注释您的存储库,但我不建议这样做,因为存储库代表数据库中对特定表的简单操作。所以存储库不知道事务上下文,因为在一个事务中可以对多个不同的存储库进行多个不同的调用。

你也可以看到我去年准备的技术演讲:https://www.youtube.com/watch?v=_1QPCoCsCTY

现在,我正在准备该主题的更新。

季博
2023-03-14

正如本文中关于在R2DBC中使用jOOQ的解释,使用R2DBC管理资源的一个好方法是使用Flux.using当(),例如。

Flux.usingWhen(
        pool.create(),
        c -> c.createStatement("SELECT col FROM my_table").execute(),
        c -> c.close()
    )
    .flatMap(it -> it.map((r, m) -> r.get(0, String.class)))
    .doOnNext(System.out::println)
    .subscribe();

邮件列表中也建议:

  • https://groups.google.com/g/r2dbc/c/a7CQAU_u_m0/m/XoRxMPXRBQAJ
  • https://groups.google.com/g/r2dbc/c/nyIXQ0EQddQ/m/HQ5rRYTnBQAJ

并有望记录在r2dbc上。未来io网站

  • https://github.com/r2dbc/r2dbc.github.io/issues/30
 类似资料:
  • 我无法使用spring webflux和r2dbc(使用r2dbc池驱动程序)打开超过10个连接。我的配置如下所示: 当我指定10个以上的连接时,会出现如下错误: 此外,连接的数量保持与初始大小相同。未创建新连接。

  • 我得到以下exeption连接到Mssql服务器。 我在属性中使用相同的配置连接到JDBC,但在尝试连接到R2DBC时出现问题。在Rest时发生,而不是在启动应用程序时发生。

  • 我正在使用macOS BigSur。我使用ssh隧道在远程机器上的gpu上运行脚本。由于进程很长,我正在使用,我希望在断开与的连接时,进程继续运行。但问题是,当ssh断开时,

  • 我正在使用spring rest模板发送与apache http client 4.2.1集成的rest请求。 由于需要向多个服务器发送请求,增加了PoolingClientConnectionManager来管理连接。 当系统运行几天后,我们发现连接达到了最大每路由设置。 打印日志如下所示保持活动的总数:0;分配路线:5选5;分配总数:100个中的5个 似乎由于某种原因,连接没有被释放。但是当我

  • 如果要从单个进程连接到数据库,则应仅创建一个 Sequelize 实例. Sequelize 将在初始化时建立连接池. 可以通过构造函数的 options 参数(使用 options.pool)来配置此连接池,如以下示例所示: const sequelize = new Sequelize(/* ... */, { // ... pool: { max: 5, min: 0

  • C3P0不会在事务完成后释放连接。下面是堆栈跟踪: 池配置和事务配置如下: 如有任何建议,我将不胜感激