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

支持超时自动更新且更新期间旧数据可用的Java/Scala集合选型咨询

实现带有定时自动更新的线程安全集合(适配Flink流场景)

当然有办法搞定这个需求!本质上我们要的是一个线程安全、支持原子切换数据视图的容器,再搭配定时任务线程定期刷新数据。结合你的Flink流数据增强场景,我分Java和Scala两种方案来给你拆解:

Java 实现方案

核心思路

用AtomicReference封装你的目标集合(比如ArrayList或HashMap,按需选择),它支持原子性的更新操作,能保证替换数据时的线程安全。再用ScheduledExecutorService启动定时任务,定期从PostgreSQL全量拉取数据,拉取完成后直接替换AtomicReference里的集合实例——这样原数据的访问完全不受影响,因为旧的集合实例还在被持有,新集合是独立的副本。

代码示例

import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;

public class AutoUpdatingCollection<T> {
    private final AtomicReference<List<T>> dataRef;
    private final ScheduledExecutorService scheduler;
    private final String dbUrl;
    private final String dbUser;
    private final String dbPassword;
    private final long updateInterval;
    private final TimeUnit timeUnit;

    // 初始化初始数据并启动定时任务
    public AutoUpdatingCollection(String dbUrl, String dbUser, String dbPassword, long updateInterval, TimeUnit timeUnit) {
        this.dbUrl = dbUrl;
        this.dbUser = dbUser;
        this.dbPassword = dbPassword;
        this.updateInterval = updateInterval;
        this.timeUnit = timeUnit;
        this.dataRef = new AtomicReference<>(fetchDataFromDB());
        this.scheduler = Executors.newSingleThreadScheduledExecutor();
        // 立即执行一次更新,之后按指定间隔循环执行
        scheduler.scheduleAtFixedRate(this::updateData, 0, updateInterval, timeUnit);
    }

    // 对外暴露的安全获取方法:返回不可变副本,防止外部修改内部状态
    public List<T> getCurrentData() {
        return List.copyOf(dataRef.get());
    }

    // 从PostgreSQL全量拉取数据的核心方法
    private List<T> fetchDataFromDB() {
        List<T> newData = new ArrayList<>();
        try (Connection conn = DriverManager.getConnection(dbUrl, dbUser, dbPassword);
             Statement stmt = conn.createStatement();
             ResultSet rs = stmt.executeQuery("SELECT * FROM your_table")) {
            // 这里替换成你的表结构转实体逻辑
            while (rs.next()) {
                T item = (T) rs.getString("target_column"); // 示例转换,按需修改
                newData.add(item);
            }
        } catch (Exception e) {
            e.printStackTrace();
            // 拉取失败时返回旧数据,避免业务中断
            return dataRef.get() != null ? dataRef.get() : newData;
        }
        return newData;
    }

    // 原子更新数据,线程安全
    private void updateData() {
        List<T> newData = fetchDataFromDB();
        dataRef.set(newData);
    }

    // 应用关闭时释放资源,避免内存泄漏
    public void shutdown() {
        scheduler.shutdown();
        try {
            if (!scheduler.awaitTermination(5, TimeUnit.SECONDS)) {
                scheduler.shutdownNow();
            }
        } catch (InterruptedException e) {
            scheduler.shutdownNow();
        }
    }
}

Flink适配注意点

  • 在Flink算子中使用时,建议在open()方法中初始化这个容器,确保每个并行算子实例只启动一个定时任务(除非你需要每个并行度独立拉取数据)。
  • 如果需要所有并行算子共享同一份更新后的数据,可以结合Flink的广播状态,把定时拉取的数据广播到所有算子实例,实现全局一致的增强数据。

Scala 实现方案

Scala可以复用Java并发包的AtomicReference,或者用Scala原生的并发工具,配合Akka调度器实现定时任务,写法更简洁。

代码示例

import java.sql.{Connection, DriverManager, ResultSet, Statement}
import java.util.concurrent.atomic.AtomicReference
import scala.concurrent.duration._
import scala.concurrent.ExecutionContext.Implicits.global

class AutoUpdatingCollection[T](dbUrl: String, dbUser: String, dbPassword: String, updateInterval: FiniteDuration) {
  private val dataRef = new AtomicReference[List[T]](fetchDataFromDB())

  // 启动Akka定时调度器
  private val scheduler = akka.actor.ActorSystem("data-updater").scheduler.scheduleAtFixedRate(
    initialDelay = 0.seconds,
    interval = updateInterval
  )(updateData)

  // 获取当前数据,返回Scala不可变List
  def getCurrentData: List[T] = dataRef.get()

  // 全量拉取PostgreSQL数据
  private def fetchDataFromDB(): List[T] = {
    val newData = scala.collection.mutable.ListBuffer[T]()
    var conn: Connection = null
    var stmt: Statement = null
    var rs: ResultSet = null
    try {
      conn = DriverManager.getConnection(dbUrl, dbUser, dbPassword)
      stmt = conn.createStatement()
      rs = stmt.executeQuery("SELECT * FROM your_table")
      while (rs.next()) {
        // 替换成你的实体转换逻辑
        val item = rs.getString("target_column").asInstanceOf[T]
        newData += item
      }
    } catch {
      case e: Exception =>
        e.printStackTrace()
        // 拉取失败返回旧数据
        return dataRef.get()
    } finally {
      if (rs != null) rs.close()
      if (stmt != null) stmt.close()
      if (conn != null) conn.close()
    }
    newData.toList
  }

  // 原子更新数据
  private def updateData(): Unit = {
    val newData = fetchDataFromDB()
    dataRef.set(newData)
  }

  // 关闭调度器和Akka系统
  def shutdown(): Unit = {
    scheduler.cancel()
    akka.actor.ActorSystem("data-updater").terminate()
  }
}

Flink适配注意点

  • 在Flink算子的close()方法中调用shutdown(),释放Akka调度器资源,避免内存泄漏。
  • 若需分布式共享数据,同样可以结合Flink广播状态实现全局数据同步。

通用最佳实践

  • 线程安全保障:永远不要让外部直接修改内部集合,返回时尽量返回不可变副本(Java的List.copyOf()、Scala的不可变List),避免并发修改问题。
  • 异常容错:数据库拉取失败时,一定要保留旧数据,不能让集合为空或抛出异常中断业务。
  • 性能优化:如果表数据量较大,拉取时用临时集合存储新数据,完成后再原子替换,避免影响原数据访问;同时添加更新日志,记录更新时间和数据量,方便排查问题。

内容的提问来源于stack exchange,提问作者L. Viktor

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:36:11