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

Vertx定时轮询中MySQL PreparedStatement后续查询返回0行问题

问题与解决方案:Vertx定时轮询数据库仅首次查询有效

问题描述

我正在移植一个每隔几分钟轮询数据库的旧Java程序,需基于表中数据处理(无法改用CDC方案,必须保留轮询机制)。使用Vertx的setTimer实现串行定时调度(确保前一次任务结束后再启动下一次),类中定义了private static Handler<Long> pollingHandler;属性并初始化,首次调用正常。将FETCH_SAMPLE_DATA_QUERY配合有效的start_id和batch_size在MySQL控制台执行可正常返回数据,但代码中仅首次调用能返回正确行数,后续所有调度调用均返回0行,日志满是Processing 0 records。调试发现首次调用后start_id已更新,但后续调用中start_id的日志值保持不变,用该调试值在控制台查询可返回数据,无法定位问题。

代码片段:

pollingHandler = id -> {
   LOGGER.debug("start_id: {}", start_id);
   LOGGER.debug("batch_size: {}", batch_size);

   mysqlClient
    .preparedQuery(FETCH_SAMPLE_DATA_QUERY)
    .execute(Tuple.of(start_id,batch_size))
    .onSuccess(rows -> {
      LOGGER.debug("Processing {} rows", rows.size());
      
      if (rows.size() > 0) {
        RowIterator<Row> rowIterator = rows.iterator();
        processRows(rowIterator); // 该方法更新start_id供下一次查询使用
      }
      vertx.setTimer(config.getLong("polling-interval"), pollingHandler);
    })
    .onFailure(err -> {
      LOGGER.error(err.getMessage(),err);
      vertx.setTimer(config.getLong("polling-interval"), pollingHandler);
    });
 };

 // 首次调度执行
 vertx.setTimer(config.getInteger("polling-interval"), pollingHandler);

问题根源

  1. 变量捕获与作用域冲突:如果start_id是初始化pollingHandler的方法内的局部变量,lambda会捕获该变量的副本(要求变量为effectively final),后续外部对start_id的修改无法同步到lambda内部,导致每次查询都使用初始值;若局部变量与类成员变量同名,lambda捕获局部变量,而processRows修改的是类成员变量,两者值完全脱节。
  2. 线程可见性问题:Vertx是多线程异步框架,基本类型的start_id在多线程更新时可能出现可见性问题,lambda无法读取到最新值。

解决方案

1. 统一变量作用域

将start_id定义为类级别的静态成员变量(与pollingHandler的static属性匹配),避免局部变量同名:

private static long start_id = 0; // 类静态变量,确保所有调用共享同一值
private static int batch_size = 100;

2. 保障线程安全与可见性

使用AtomicLong替代基本类型long,确保多线程环境下的原子更新与可见性:

private static AtomicLong start_id = new AtomicLong(0);

修改后的核心代码:

pollingHandler = id -> {
    long currentStartId = start_id.get(); // 获取最新值
    LOGGER.debug("start_id: {}", currentStartId);
    LOGGER.debug("batch_size: {}", batch_size);

    mysqlClient
        .preparedQuery(FETCH_SAMPLE_DATA_QUERY)
        .execute(Tuple.of(currentStartId, batch_size))
        .onSuccess(rows -> {
            LOGGER.debug("Processing {} rows", rows.size());
            
            if (rows.size() > 0) {
                RowIterator<Row> rowIterator = rows.iterator();
                // 假设processRows返回新的start_id值
                long newStartId = processRows(rowIterator);
                start_id.set(newStartId); // 原子更新值
            }
            vertx.setTimer(config.getLong("polling-interval"), pollingHandler);
        })
        .onFailure(err -> {
            LOGGER.error(err.getMessage(), err);
            vertx.setTimer(config.getLong("polling-interval"), pollingHandler);
        });
};

3. 验证processRows逻辑

确保processRows确实在修改正确的start_id变量:

  • 若使用AtomicLong,需调用set()或getAndSet()方法更新;
  • 若使用静态long变量,可添加volatile修饰符保障可见性,直接赋值即可。

4. 修正调度参数类型

首次调度时统一参数类型,将config.getInteger("polling-interval")改为config.getLong("polling-interval"),避免类型转换导致的异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 06:00:07