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

Vert.x应用中多批次运行的两个Verticle出现MySQL超时问题求助

Vert.x应用中多批次运行的两个Verticle出现MySQL超时问题求助

大家好,我在开发Vert.x应用时遇到了MySQL相关的超时问题,目前有两个负责多批次任务的Verticle(TestVerticle4和TestVerticle7),运行过程中总会出现数据库连接超时或者操作超时的情况,想请各位帮忙排查下问题所在。

先贴出两个Verticle的代码:

TestVerticle4 代码

public class TestVerticle4 extends AbstractVerticle{
    private static final Logger logger = LoggerFactory.getLogger(TestVerticle4.class);

    @Override
    public void start(Promise<Void> startPromise) {
        try {
            MySQLConnectOptions connectOptions = new MySQLConnectOptions()
                    .setPort(Integer.parseInt(config().getString("sqlPort")))
                    .setHost(config().getString("host"))
                    .setDatabase(config().getString("database"))
                    .setUser(config().getString("user"))
                    .setPassword(config().getString("password"))
                    .setConnectTimeout(30000)
                    .setIdleTimeout(300)
                    .setConnectTimeout(5000)
                    .setReconnectInterval(2000);

            PoolOptions poolOptions = new PoolOptions().setMaxSize(60);
            MySQLPool client = MySQLPool.pool(vertx,connectOptions,poolOptions);

            CircuitBreakerOptions cbOptions = new CircuitBreakerOptions()
                    .setMaxRetries(3) // 重试3次
                    .setResetTimeout(5000) // 成功后重置超时时间(毫秒)
                    .setTimeout(1000) // 操作超时时间(毫秒)
                    .setFallbackOnFailure(true); // 开启失败降级
            CircuitBreaker breaker = CircuitBreaker.create("TestVerticle7", vertx, cbOptions);

            logger.info("Step3:------------ Recieving at Verticle4 ");
            logger.info("Recieving at task.runVerticle4....");

            EventBus bus =vertx.eventBus();
            bus.<JsonObject>consumer("task.runVerticle4",msg->{
                JsonObject bodyCon = msg.body();
                Long hashcode=bodyCon.getLong("hashcode");
                Long minId= bodyCon.getLong("minId");
                Long maxId= bodyCon.getLong("maxId");
                Long prodId=bodyCon.getLong("prodId");

                breaker.execute(promise -> {
                    client.getConnection( ar->{
                        if(ar.succeeded()) {
                            SqlConnection conn =ar.result();
                            conn.preparedQuery(insertQueryEvent())
                                    .execute(Tuple.of(hashcode,prodId,minId,maxId),res->{
                                        if(res.succeeded()) {
                                            logger.info("Successfully Inserted Batch for custIds {} to {} for product {} for hashcode {}",minId,maxId,prodId,hashcode);
                                        }else {
                                            logger.error("Failed to execute "+res.cause().getMessage());
                                        }
                                        conn.close();
                                    });
                        }else {
                            logger.error("Failed to connect to db"+ ar.cause().getMessage());
                        }
                    });
                // 原代码此处的onComplete回调被注释,是否为故意操作?
                // }).onComplete(ar -> {
                //     if (ar.succeeded()) {
                //         System.out.println("Query executed successfully (or fallback used)");
                //     } else {
                //         System.err.println("Query failed after retries: " + ar.cause().getMessage());
                //     }
                // });
                });
            });
            startPromise.complete();
        }catch(Exception e) {
            e.printStackTrace();
            startPromise.fail(e);
        }
    }

    public static String insertQueryEvent() {
        return " INSERT IGNORE INTO ResearchCall_HashQueue(Customer_Id,Product_Id,HashCode) " +
                "SELECT Customer_Id,Reco_Product_1, ? FROM RankedPickList_Stories " +
                "WHERE Reco_Product_1= ? AND Customer_Id BETWEEN ? AND ? AND Story_Id=3;";
    }
}

TestVerticle7 代码(注:原代码存在截断,以下为现有完整部分)

public class TestVerticle7 extends AbstractVerticle{
    private static final Logger logger = LoggerFactory.getLogger(TestVerticle7.class);

    @Override
    public void start(Promise<Void> startPromise) {
        try {
            MySQLConnectOptions connectOptions = new MySQLConnectOptions()
                    .setPort(Integer.parseInt(config().getString("sqlPort")))
                    .setHost(config().getString("host"))
                    .setDatabase(config().getString("database"))
                    .setUser(config().getString("user"))
                    .setPassword(config().getString("password"))
                    .setConnectTimeout(30000)
                    .setIdleTimeout(300)
                    .setConnectTimeout(5000)
                    .setReconnectInterval(2000);

            PoolOptions poolOptions = new PoolOptions().setMaxSize(60);
            MySQLPool client = MySQLPool.pool(vertx,connectOptions,poolOptions);

            CircuitBreakerOptions cbOptions = new CircuitBreakerOptions()
                    .setMaxRetries(3) // 重试3次
                    .setResetTimeout(5000) // 成功后重置超时时间(毫秒)
                    .setTimeout(1000) // 操作超时时间(毫秒)
                    .setFallbackOnFailure(true); // 开启失败降级
            CircuitBreaker breaker = CircuitBreaker.create("TestVerticle7", vertx, cbOptions);

            logger.info("Recieving to task.runVerticle7....");
            EventBus bus =vertx.eventBus();
            bus.<JsonObject>consumer("task.runVerticle7",msg->{
                JsonObject msg_body = msg.body();
                Long minId= msg_body.getLong("minId");
                Long maxId= msg_body.getLong("maxId");
                Long prodId= msg_body.getLong("prodId");
                Long hashcode= msg_body.getLong("hashcode");

                breaker.execute(promise -> {
                    client.getConnection( ar->{
                        if(ar.succeeded()) {
                            SqlConnection conn =ar.result();
                            conn.preparedQuery(updateProductArray())
                                    .execute(Tuple.of(minId,maxId,prodId,hashcode),res->{
                                        if(res.succeeded()) {
                                            logger.info("Successfull Updated product array to Customer_ResearchCall_Master for Cust_Ids between {} and {}",minId,maxId);
                                        }else {
                                            logger.error("Failed to update product array to Customer_Res..."); // 原代码此处截断
                                        }
                                        // 此处疑似忘记关闭SqlConnection?
                                    });
                        }else {
                            logger.error("Failed to connect to db"+ ar.cause().getMessage());
                        }
                    });
                });
            });
            startPromise.complete();
        }catch(Exception e) {
            e.printStackTrace();
            startPromise.fail(e);
        }
    }

    // 原代码中updateProductArray方法未给出完整实现
    public static String updateProductArray() {
        return "";
    }
}

我目前观察到几个可疑点,想和大家确认:

  1. 连接超时设置冲突:在MySQLConnectOptions里重复调用了setConnectTimeout,先设30000ms又设5000ms,实际生效的是最后一个5000ms,这个时长会不会过短?
  2. 断路器与数据库超时不匹配:断路器超时设置为1000ms,但数据库连接超时是5000ms,会不会导致断路器先触发超时,而数据库连接还在尝试建立?
  3. 资源释放遗漏:TestVerticle7中获取SqlConnection后,好像忘记调用conn.close(),会不会导致连接池被耗尽,进而引发后续连接超时?
  4. 连接池负载问题:两个Verticle的连接池都设置了maxSize=60,多批次任务同时运行时,会不会超过数据库实例允许的最大连接数?
  5. 断路器名称错误:TestVerticle4里创建的断路器名称是"TestVerticle7",是不是写错了?会不会导致两个Verticle共用同一个断路器实例引发冲突?

运行时频繁收到数据库连接超时或操作超时的错误日志,想请各位帮忙看看代码里还有哪些潜在问题,或者有没有优化方向?

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 11:18:01