如何将HTTP API封装为Flink Source并获取Java实现示例
基于Java实现接收HTTP请求的Flink Source
实现思路
- 自定义Flink Source内部启动嵌入式HTTP服务器,监听指定端口接收请求
- HTTP请求到达后先向调用方返回接收成功响应,再将请求体数据写入内存阻塞队列
- Source消费阻塞队列中的数据,下发给后续Flink算子处理
依赖配置
首先在pom.xml中引入所需依赖,除了基础Flink依赖外,额外引入Jetty嵌入式服务器组件:
<dependencies> <!-- Flink 基础依赖,根据你的集群版本调整 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.1</version> <scope>provided</scope> </dependency> <!-- Jetty 嵌入式HTTP服务 --> <dependency> <groupId>org.eclipse.jetty</groupId> <artifactId>jetty-server</artifactId> <version>9.4.48.v20220622</version> </dependency> <dependency> <groupId>org.eclipse.jetty</groupId> <artifactId>jetty-servlet</artifactId> <version>9.4.48.v20220622</version> </dependency> </dependencies>
完整代码实现
自定义HTTP Source类
import org.apache.flink.streaming.api.functions.source.SourceFunction; import org.eclipse.jetty.server.Server; import org.eclipse.jetty.servlet.ServletContextHandler; import org.eclipse.jetty.servlet.ServletHolder; import javax.servlet.http.HttpServlet; import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; import java.io.BufferedReader; import java.io.IOException; import java.io.InputStreamReader; import java.nio.charset.StandardCharsets; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; public class HttpSource implements SourceFunction<String> { private final int port; private transient Server server; private transient BlockingQueue<String> dataQueue; private volatile boolean isRunning = true; public HttpSource(int port) { this.port = port; } @Override public void run(SourceContext<String> ctx) throws Exception { dataQueue = new LinkedBlockingQueue<>(1000); // 队列容量可根据业务调整 // 启动Jetty服务 server = new Server(port); ServletContextHandler context = new ServletContextHandler(ServletContextHandler.SESSIONS); context.setContextPath("/"); server.setHandler(context); // 注册接收请求的Servlet context.addServlet(new ServletHolder(new HttpServlet() { @Override protected void doPost(HttpServletRequest req, HttpServletResponse resp) throws IOException { // 1. 先返回响应给调用方 resp.setStatus(HttpServletResponse.SC_OK); resp.setContentType("application/json"); resp.getWriter().write("{\"code\":0,\"message\":\"数据接收成功\"}"); resp.getWriter().flush(); // 2. 读取请求体写入队列 BufferedReader reader = new BufferedReader(new InputStreamReader(req.getInputStream(), StandardCharsets.UTF_8)); StringBuilder sb = new StringBuilder(); String line; while ((line = reader.readLine()) != null) { sb.append(line); } dataQueue.offer(sb.toString()); } }), "/receive"); // 接口路径可自定义 server.start(); // 循环消费队列数据下发给下游 while (isRunning) { String data = dataQueue.take(); synchronized (ctx.getCheckpointLock()) { ctx.collect(data); } } } @Override public void cancel() { isRunning = false; try { if (server != null) { server.stop(); } } catch (Exception e) { e.printStackTrace(); } if (dataQueue != null) { dataQueue.clear(); } } }
作业使用示例
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class HttpSourceJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 设置Source并行度为1,避免多实例端口冲突 env.addSource(new HttpSource(8888)) .setParallelism(1) .print(); // 此处替换为你的业务处理逻辑 env.execute("HTTP Source Demo"); } }
使用说明
- 作业启动后,向
http://<Flink TaskManager地址>:8888/receive发送POST请求,即可自动将请求体传入Flink流处理链路 - 若需要高并发接收请求,可将Source并行度设置为N,同时配置N个不同端口,在前端部署负载均衡将请求分发到各个端口
- 若需要实现精确一次语义,可结合Flink Checkpoint机制,将未下发的数据存储到ListState中,故障恢复时从State中恢复数据
内容的提问来源于stack exchange,提问作者Gabriel Plutov
相关产品推荐
相关产品推荐

