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

如何在Flink自定义Source中使用ScheduledExecutor实现每小时HTTP轮询

Flink定时轮询HTTP接口自定义Source实现说明

可行性确认

可以在自定义Source内部使用ScheduledExecutor实现固定间隔轮询逻辑,你现有代码的问题是只会执行一次HTTP请求,执行完成后run方法退出,Source会直接终止运行,无法持续轮询。

改造后的代码实现

你可以按照如下逻辑修改自定义Source:

import scala.io.Source.fromInputStream
import com.typesafe.scalalogging.LazyLogging
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.source.{RichSourceFunction, SourceFunction}
import java.util.concurrent.{Executors, ScheduledExecutorService, TimeUnit}

class HttpSource(url: String, pollIntervalHours: Int) extends RichSourceFunction[String] with LazyLogging {

  @volatile private var isRunning = true
  private var executor: ScheduledExecutorService = _

  override def open(parameters: Configuration): Unit = {
    // 初始化单线程调度器,避免多线程重复拉取
    executor = Executors.newSingleThreadScheduledExecutor()
  }

  override def cancel(): Unit = {
    isRunning = false
  }

  override def run(ctx: SourceFunction.SourceContext[String]): Unit = {
    // 提交定时轮询任务,首次立即执行,之后按指定间隔执行
    executor.scheduleAtFixedRate(() => {
      if (isRunning) {
        httpStream(ctx.collect)
      }
    }, 0, pollIntervalHours, TimeUnit.HOURS)

    // 阻塞run方法,直到作业停止
    while (isRunning) {
      Thread.sleep(1000)
    }
  }

  private def httpStream(rec: String => Unit): Unit = {
    try {
      val request = Http(url)
      val response = request.execute()
      if (response.code != 200) {
        logger.warn("Metadata api response status code: {}", response.code)
      } else {
        request.execute { inputStream =>
          fromInputStream(inputStream)
            .getLines()
            .foreach(rec)
        }
      }
    } catch {
      case e: Exception => 
        logger.error("Poll http api failed", e)
    }
  }

  override def close(): Unit = {
    // 作业停止时主动关闭线程池,避免资源泄漏
    if (executor != null) {
      executor.shutdownNow()
    }
  }
}

作业使用方式

// 每1小时轮询一次接口
val httpSource = env.addSource(new HttpSource(baseUri, 1))
httpSource.broadcast.print()

关键注意事项

  • 必须保证isRunning变量加@volatile注解,确保多线程之间的状态可见性,避免cancel操作不生效
  • 定时任务内部要捕获所有异常,不要抛出未处理的异常,否则会导致调度线程终止,后续轮询任务不会再执行
  • 不要在run方法里直接返回,需要加阻塞逻辑保持run方法持续运行,否则Source会被标记为结束,作业会停止
  • 如果要实现 Exactly-Once 语义,可以配合Flink的状态机制,将上次拉取的时间、偏移量等信息存入托管状态,作业故障恢复后可以从对应位置继续拉取,避免重复或者丢失数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 00:27:01