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

如何将HTTP API封装为Flink Source并获取Java实现示例

实现思路

  • 自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 08:18:01