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

虚拟线程、HikariCP与Spring Boot JDBC(Postgres)并发IO问题排查

高并发下HikariCP搭配Project Loom出现PostgreSQL IO错误问题

技术栈

  • Docker容器运行的PostgreSQL
  • Project Loom
  • Spring Boot 3.3.4
  • Java 21
  • Spring Boot Starter JDBC(内置HikariCP,疑似问题根源)
  • Spring Boot Starter JPA

所有组件(含PostgreSQL驱动)均使用最新版本。

核心问题

高并发操作场景下,通过HikariCP向PostgreSQL写入数据时出现IO错误:Hikari标记数据库连接为损坏,PostgreSQL端日志持续显示客户端连接意外EOF且存在未关闭的事务。

日志信息

微服务日志

2024-09-19T18:50:51.617-04:00  WARN 10060 --- [virtual-ingest] [    virtual-342] com.zaxxer.hikari.pool.ProxyConnection   : HikariPool-1 - Connection org.postgresql.jdbc.PgConnection@28b26c3f marked as broken because of SQLSTATE(08006), ErrorCode(0)

org.postgresql.util.PSQLException: An I/O error occurred while sending to the backend.

PostgreSQL Docker日志(重复出现)

2024-09-19 22:50:51.807 UTC [412] LOG:  unexpected EOF on client connection with an open transaction

相关代码

消息处理结构代码

@EventListener(ApplicationReadyEvent.class)
public void startStructures()
{
    try (var outerScope = new StructuredTaskScope.ShutdownOnFailure())
    {
        sqsConfig.getQueues().forEach(queue -> outerScope.fork(() -> {
            processQueueMessages(queue);
            return null;
        }));
        outerScope.join();
    }
    catch (InterruptedException e)
    {
        Thread.currentThread().interrupt();
        log.error("Processing interrupted", e);
    }
}

private void processQueueMessages(String queue)
{
    while (true)
    {
        try (var innerScope = new StructuredTaskScope.ShutdownOnFailure())
        {
            List<Message> messages = retrieveMessages(queue);
            if (messages.isEmpty())
            {
                log.info("No messages in queue: {}, sleeping...", queue);
                Thread.sleep(5000);
            }
            else
            {
                for (Message message : messages)
                {
                    innerScope.fork(() -> {
                        processMessage(queue, message);
                        return null;
                    });
                }
                innerScope.join();
            }
        }
        catch (InterruptedException e)
        {
            Thread.currentThread().interrupt();
            log.error("Queue processing interrupted", e);
            break;
        }
    }
}

public List<Message> retrieveMessages(String queue)
{
    var url = awsConfig.getBaseUrl() + queue;
    ReceiveMessageRequest request = ReceiveMessageRequest.builder()
                                                         .queueUrl(url)
                                                         .maxNumberOfMessages(10)
                                                         .waitTimeSeconds(10)
                                                         .build();

    ReceiveMessageResponse response = sqsClient.receiveMessage(request);
    return response.messages();
}

private void processMessage(String queue, Message message) throws JsonProcessingException
{
    log.info("Processing message from queue {}: {}", queue, message.messageId());

    JsonNode node = nodeBuilderService.buildNode(message);
    distributionService.handleSqsNotification(node);
    deleteMessage(queue, message);
}

private void deleteMessage(String queue, Message message)
{
    var url = awsConfig.getBaseUrl() + queue;
    sqsClient.deleteMessage(builder -> builder.queueUrl(url).receiptHandle(message.receiptHandle()));

    int remainingMessages = getMessageCount(queue);
    log.info(
        "Deleted message: {}. Approximate {} messages remaining in the queue.", message.messageId(),
        remainingMessages
    );
}

public int getMessageCount(String queueName)
{
    var url = awsConfig.getBaseUrl() + queueName;

    GetQueueAttributesRequest request = GetQueueAttributesRequest.builder()
                                                                 .queueUrl(url)
                                                                 .attributeNames(QueueAttributeName.APPROXIMATE_NUMBER_OF_MESSAGES)
                                                                 .build();

    return Integer.parseInt(
        sqsClient.getQueueAttributes(request)
                 .attributes()
                 .get(QueueAttributeName.APPROXIMATE_NUMBER_OF_MESSAGES)
    );
}

触发错误的事务方法

@Transactional
public void process(String filename, InputStream stream) throws IOException
{
    byte[] bytes = getBytesFromInputStream(stream, filename);

    List<SurfaceObservation> observations = decoder.beginSynopticDecoders(bytes, filename);
    repository.saveAllAndFlush(observations);
}

疑问

  1. 上述代码是否存在明显错误?正从异步模式迁移至基于StructuredTaskScope的实现,对该API的正确用法可能存在误解。
  2. ExtentLocal能否解决此问题?怀疑HikariCP依赖ThreadLocal的实现与虚拟线程不兼容,而ExtentLocal是ThreadLocal的替代方案。
  3. 是否必须放弃Spring JDBC?若需替换,有哪些可行方案?目前认为R2DBC(搭配Hibernate Reactive)的ORM生态仍不够成熟。

已尝试的操作:强制使用平台线程、为数据库事务单独创建StructuredTaskScope等,但未能解决问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 04:57:32