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

如何实现每小时轮询HTTP端点的Flink流SourceFunction数据源

现有代码核心问题

  • 需求是主动定时拉取外部HTTP接口,但当前代码实现的是启动本地HTTP服务端,被动等待外部请求调用,逻辑完全相反
  • server.awaitTermination(1, TimeUnit.HOURS) 只会让服务运行1小时就自动终止,不符合流作业长期运行的要求
  • 没有实现每小时固定间隔触发请求的逻辑

正确实现方案

改用HTTP客户端发起请求,配合循环休眠实现每小时轮询逻辑,示例代码如下:

自定义轮询数据源实现

import org.apache.flink.api.common.typeinfo.BasicTypeInfo
import org.apache.flink.streaming.api.functions.source.SourceFunction
import org.apache.flink.streaming.api.functions.source.SourceFunction.SourceContext
import org.apache.http.client.methods.HttpGet
import org.apache.http.impl.client.HttpClients
import org.apache.http.util.EntityUtils
import java.util.concurrent.TimeUnit

class HourlyHttpPollingSource(url: String) extends SourceFunction[String] {
  // 作业运行标记,cancel时终止循环
  @volatile private var isRunning = true
  // HTTP客户端不可序列化,加transient注解
  @transient private var httpClient = _

  override def run(ctx: SourceContext[String]): Unit = {
    httpClient = HttpClients.createDefault()
    // 如果需要启动时立即执行一次请求,放开下面的注释
    // pollOnce(ctx)
    while (isRunning) {
      // 每小时轮询一次
      TimeUnit.HOURS.sleep(1)
      if (isRunning) {
        pollOnce(ctx)
      }
    }
  }

  private def pollOnce(ctx: SourceContext[String]): Unit = {
    val get = new HttpGet(url)
    try {
      val response = httpClient.execute(get)
      val statusCode = response.getStatusLine.getStatusCode
      if (statusCode == 200) {
        val responseContent = EntityUtils.toString(response.getEntity, "UTF-8")
        // 加锁保证数据收集和快照操作的原子性
        ctx.getCheckpointLock.synchronized {
          ctx.collect(responseContent)
        }
      }
      EntityUtils.consume(response.getEntity)
    } catch {
      case e: Exception =>
        // 可根据需要添加重试、告警逻辑,避免单次请求失败导致作业终止
        e.printStackTrace()
    } finally {
      get.releaseConnection()
    }
  }

  override def cancel(): Unit = {
    isRunning = false
    if (httpClient != null) {
      httpClient.close()
    }
  }
}

主作业广播流配置

// 数据源并行度设置为1,避免多实例重复发起轮询请求
val pollingStream = env.addSource(new HourlyHttpPollingSource(config.baseUri)).setParallelism(1)

// 定义广播状态描述符,可根据业务需要调整存储的KV类型
val broadcastDescriptor = new MapStateDescriptor[String, String](
  "httpPollingBroadcast",
  BasicTypeInfo.STRING_TYPE_INFO,
  BasicTypeInfo.STRING_TYPE_INFO
)
val broadcastStream = pollingStream.broadcast(broadcastDescriptor)

// 下游普通流连接广播流即可获取到广播的轮询数据
// 示例:普通数据流关联广播流处理
// val resultStream = commonStream.connect(broadcastStream).process(yourCustomBroadcastProcessFunction)

// 测试打印
pollingStream.print()

注意事项

  • 如果需要实现Exactly Once语义,可以让数据源实现CheckpointedFunction接口,在快照中存储上次轮询的时间戳,作业故障恢复后从上次状态继续执行
  • 可以根据业务需要调整HTTP请求的超时时间、重试次数、请求头等参数

内容的提问来源于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 13:54:03