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

Apache Flink JDBC Source定时执行与持续运行问题咨询

默认的JdbcSource是为批处理场景设计的,执行完一次查询就会终止作业。要实现每N秒调度执行并让作业持续运行,推荐以下两种方案:

方案一:自定义周期性JDBC Source(最简实现)

自己实现一个RichParallelSourceFunction,在循环中周期性执行JDBC查询,同时处理作业取消与异常重连逻辑:

class PeriodicJdbcSource(
    private val dbUrl: String,
    private val username: String,
    private val password: String,
    private val querySql: String,
    private val intervalSeconds: Long,
    private val resultExtractor: (ResultSet) -> LoggedInEvent
) : RichParallelSourceFunction<LoggedInEvent>() {

    private var isRunning = true
    private var connection: Connection? = null
    private var statement: Statement? = null

    override fun open(parameters: Configuration?) {
        // 加载驱动并初始化JDBC连接
        Class.forName("org.postgresql.Driver")
        connection = DriverManager.getConnection(dbUrl, username, password)
        statement = connection?.createStatement()
    }

    override fun run(ctx: SourceContext<LoggedInEvent>) {
        while (isRunning) {
            try {
                // 执行查询并发送数据到流中
                statement?.executeQuery(querySql)?.use { resultSet ->
                    while (resultSet.next()) {
                        val event = resultExtractor(resultSet)
                        ctx.collect(event)
                    }
                }
                // 等待指定间隔后再次执行
                Thread.sleep(intervalSeconds * 1000)
            } catch (e: InterruptedException) {
                // 捕获中断信号,终止循环
                isRunning = false
            } catch (e: SQLException) {
                // 处理JDBC异常,比如断开后重连
                e.printStackTrace()
                closeResources()
                open(null)
            }
        }
    }

    override fun cancel() {
        isRunning = false
        closeResources()
    }

    private fun closeResources() {
        statement?.runCatching { close() }
        connection?.runCatching { close() }
    }
}

使用方式

替换原有的JdbcSource,指定调度间隔(示例为10秒):

val periodicSource = PeriodicJdbcSource(
    dbUrl = "jdbc:postgresql://db:5432/postgres",
    username = "postgres",
    password = "example",
    querySql = "SELECT player_id, past_logins FROM user_initial_data",
    intervalSeconds = 10,
    resultExtractor = { LoggedInEvent(it.getInt(1).toString(), it.getInt(2), Instant.now().toEpochMilli()) }
)

// 添加自定义Source到环境,作业会持续运行
val snapshotsStream = env.addSource(periodicSource, "PeriodicLoggedInSnapshots")

如果业务更适合用SQL实现,可以通过注册JDBC表+窗口查询实现周期性拉取:

1. 注册JDBC表

CREATE TABLE user_initial_data (
    player_id INT,
    past_logins INT
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:postgresql://db:5432/postgres',
    'table-name' = 'user_initial_data',
    'username' = 'postgres',
    'password' = 'example'
);

2. 编写定时查询

使用TUMBLE窗口实现每N秒执行一次查询:

SELECT 
    player_id, 
    past_logins, 
    UNIX_TIMESTAMP() * 1000 AS event_time
FROM user_initial_data
TUMBLE(INTERVAL '10' SECOND); -- 每10秒执行一次

3. 转换为DataStream(可选)

如果需要用DataStream API处理,可将SQL查询结果转换为流:

val tableEnv = StreamTableEnvironment.create(env)
// 注册表...
val resultTable = tableEnv.sqlQuery("上述SQL语句")
val snapshotsStream = tableEnv.toDataStream(resultTable, LoggedInEvent::class.java)

关键说明

  • 自定义Source默认使用环境并行度,若不想重复拉取数据,可设置periodicSource.setParallelism(1)
  • 若需要更精确的定时(避免sleep累积误差),可以用ScheduledExecutorService替代Thread.sleep
  • SQL方案依赖Flink JDBC连接器版本,请确保与Flink版本兼容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 12:30:57