如何实现每小时轮询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
相关产品推荐
相关产品推荐

