如何在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
相关产品推荐
相关产品推荐

