java - r2dbc-pool 连接取消后未释放

标签 java spring spring-data-r2dbc r2dbc r2dbc-postgresql

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

  • 即使它们被取消,一些线程也会获取连接并执行正常的数据库进程
  • 最重要的是:即使所有线程都被取消且没有线程再处于 Activity 状态,某些连接仍会被获取并且永远不会被释放。

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

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

我编写了一个程序来重现该问题。

为了进行测试,我使用了 maxConnection=2 的池。 调用测试方法几次 (controller.test) 后,池中的某些连接仍然无限期地获取(它们应该全部由 onCancelclose 释放> 由 Spring 处理的语句)。通过使用 jmx 监视池可以轻松地证明这一点。

我认为取消请求传播到connectionPool.create(),但某些迭代似乎有足够的时间在收到取消之前结束预查询,这导致连接可用于 Spring用来。在这些情况下,TestConnectionFactory 中看不到取消,并且大约 1/3 的情况下,Spring 不会调用connection.close,从而导致仍然获取连接。

@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.boot 2.3.5.RELEASEio.r2dbc:r2dbc-postgresqlio.r2dbc:r2dbc-pool一起使用强>.

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

最佳答案

this article about using jOOQ with R2DBC 中所述,使用 R2DBC 管理资源的一个好方法是使用 Flux.usingWhen() ,例如

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();

邮件列表上也推荐了这一点:

并有望记录在 r2dbc.io 上 future 的网站:

关于java - r2dbc-pool 连接取消后未释放,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/68407202/

相关文章:

java - 如何让我的框架创建用户定义子类的实例

Java 内存数据库类对象

java - 在文本文档中搜索一行 - JAVA

java - 在http请求连接器中传递动态查询参数和路径参数

spring-webflux - 连接重用和事务(关系)

spring-webflux - 使用 R2DBC 的动机是什么?

SSLSocket 创建时的 Java 异常

spring - Grails:错误org.springframework.boot.SpringApplication-应用启动失败

java - Spring boot 如何选择外部化的 spring 属性文件

java - 如何在 Spring 数据 r2dbc 中替换 @PrePersist