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

如何在Flink中使用Async I/O调用分页HTTP API并实现审计追踪

基于Flink处理分页API并实现审计追踪的方案

针对你提出的分页API调用和审计追踪需求,以下是具体的Flink实现方案:

1. 批处理模式实现

批处理为有界场景,核心是迭代获取分页数据直到API返回的next指针为空,通过带状态的AsyncFunction结合迭代逻辑实现:

  • 初始触发:从配置或审计数据库读取初始next指针(首次执行可设为null),作为作业输入源(如用fromElements生成触发信号)。
  • 自定义AsyncFunction逻辑:
    1. 在asyncInvoke中,根据当前next指针构造API请求,调用HTTP客户端获取响应。
    2. 解析响应中的业务数据和新的next指针。
    3. 将业务数据输出到下游,同时将新next指针存入Flink的ValueState。
  • 迭代控制:通过Flink批处理的iterate()方法循环触发请求,直到状态中的next指针为空时终止迭代。

代码示例(简化版):

public class PagingAsyncFunction extends RichAsyncFunction<String, Tuple2<List<DataRecord>, String>> {
    private transient ValueState<String> currentNextState;
    private CloseableHttpClient httpClient;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        httpClient = HttpClients.createDefault();
        // 初始化状态,加载初始next指针
        ValueStateDescriptor<String> stateDesc = new ValueStateDescriptor<>("currentNext", String.class);
        currentNextState = getRuntimeContext().getState(stateDesc);
        currentNextState.update(parameters.getString("initial.next", null));
    }

    @Override
    public void asyncInvoke(String trigger, ResultFuture<Tuple2<List<DataRecord>, String>> resultFuture) throws Exception {
        String next = currentNextState.value();
        String apiUrl = "http://someurl.com" + (next == null ? "" : "?next=" + next);
        
        // 发送异步HTTP请求
        HttpGet request = new HttpGet(apiUrl);
        httpClient.execute(request, new AsyncResponseHandler<Void>() {
            @Override
            public Void handleResponse(HttpResponse response) throws IOException {
                List<DataRecord> dataList = parseResponseData(response.getEntity().getContent());
                String newNext = extractNextPointer(response);
                
                // 更新状态中的next指针
                currentNextState.update(newNext);
                // 输出业务数据和新next指针
                resultFuture.complete(Collections.singletonList(Tuple2.of(dataList, newNext)));
                return null;
            }
        });
    }

    @Override
    public void close() throws Exception {
        httpClient.close();
        super.close();
    }
}

// 批处理作业提交
public static void main(String[] args) throws Exception {
    ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
    env.setParallelism(1); // 分页API建议单并行度,避免重复请求

    DataSet<String> initialTrigger = env.fromElements("start");
    // 迭代处理分页数据
    DataSet<DataRecord> finalData = initialTrigger.iterate(100) {
        DataSet<String> loopInput -> {
            DataSet<Tuple2<List<DataRecord>, String>> asyncResult = AsyncDataStream.unorderedWait(
                    loopInput, new PagingAsyncFunction(), 1000, TimeUnit.MILLISECONDS, 10);
            // 输出业务数据,同时返回非空的next指针作为下一轮迭代输入
            DataSet<String> nextTrigger = asyncResult.filter(t -> t.f1 != null).map(t -> "next");
            asyncResult.flatMap((t, out) -> t.f0.forEach(out::collect)).output(new LocalOutputFormat<>());
            return nextTrigger;
        }
    };

    env.execute("Batch Paging API Job");
}

2. 流处理模式实现

流处理为无界场景,需持续获取分页数据(如定时轮询),通过状态持久化保存当前next指针:

  • 触发方式:自定义SourceFunction生成定时触发事件(如每隔5分钟触发一次)。
  • 带状态的AsyncFunction逻辑:
    1. 初始化时从审计数据库读取上次保存的next指针,存入ValueState。
    2. 每次触发请求时,从状态中取出next指针构造API请求。
    3. 解析响应后更新状态中的next指针,同时输出业务数据到下游。
  • 故障恢复:开启Checkpoint,确保next指针在作业重启时不丢失。

代码示例(简化版):

public class TimedPagingSource implements SourceFunction<String> {
    private volatile boolean running = true;
    private final long pollIntervalMs;

    public TimedPagingSource(long pollIntervalMs) {
        this.pollIntervalMs = pollIntervalMs;
    }

    @Override
    public void run(SourceContext<String> ctx) throws Exception {
        while (running) {
            ctx.collect("trigger");
            Thread.sleep(pollIntervalMs);
        }
    }

    @Override
    public void cancel() {
        running = false;
    }
}

// 流处理作业提交
public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    env.enableCheckpointing(30000); // 每30秒做一次Checkpoint

    DataStream<String> triggerStream = env.addSource(new TimedPagingSource(300000)); // 5分钟轮询一次

    AsyncDataStream.unorderedWait(triggerStream, new PagingAsyncFunction(), 1000, TimeUnit.MILLISECONDS, 10)
            .flatMap((Tuple2<List<DataRecord>, String> t, Collector<DataRecord> out) -> t.f0.forEach(out::collect))
            .print();

    env.execute("Streaming Paging API Job");
}

二、将next指针存入数据库实现审计追踪

Flink完全支持该需求,可通过JDBC Sink或Table API实现,两种方式均可将next指针与审计元数据(如处理时间、请求状态)写入数据库:

1. DataStream JDBC Sink实现

在AsyncFunction输出结果后,将next指针封装为审计记录,通过JDBC Sink写入数据库:

// 审计记录实体
public class AuditRecord {
    private String nextPointer;
    private Timestamp processTime;
    private String status; // success/failure/completed

    // 构造方法、getter/setter省略
}

// 配置JDBC Sink
JdbcExecutionOptions execOpts = JdbcExecutionOptions.builder()
        .withBatchSize(1)
        .withBatchIntervalMs(0)
        .build();

JdbcConnectionOptions connOpts = JdbcConnectionOptions.builder()
        .withUrl("jdbc:mysql://localhost:3306/audit_db")
        .withDriverName("com.mysql.cj.jdbc.Driver")
        .withUsername("db_user")
        .withPassword("db_pwd")
        .build();

// 生成审计流并写入数据库
DataStream<AuditRecord> auditStream = asyncResultStream.map(t -> {
    AuditRecord audit = new AuditRecord();
    audit.setNextPointer(t.f1);
    audit.setProcessTime(new Timestamp(System.currentTimeMillis()));
    audit.setStatus(t.f1 == null ? "completed" : "success");
    return audit;
});

auditStream.addSink(JdbcSink.sink(
        "INSERT INTO audit_table (next_pointer, process_time, status) VALUES (?, ?, ?)",
        (ps, record) -> {
            ps.setString(1, record.getNextPointer());
            ps.setTimestamp(2, record.getProcessTime());
            ps.setString(3, record.getStatus());
        },
        execOpts,
        connOpts
));

2. Table API实现

将审计流转换为Flink Table,通过INSERT INTO语句写入数据库表:

StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

// 将审计流注册为临时视图
tableEnv.createTemporaryView("audit_view", auditStream, 
        $("nextPointer"), $("processTime"), $("status"));

// 创建数据库表的映射
tableEnv.executeSql("CREATE TABLE audit_table (" +
        "next_pointer STRING," +
        "process_time TIMESTAMP," +
        "status STRING," +
        "PRIMARY KEY(next_pointer) NOT ENFORCED" +
        ") WITH (" +
        "'connector' = 'jdbc'," +
        "'url' = 'jdbc:mysql://localhost:3306/audit_db'," +
        "'table-name' = 'audit_table'," +
        "'driver' = 'com.mysql.cj.jdbc.Driver'," +
        "'username' = 'db_user'," +
        "'password' = 'db_pwd'" +
        ")");

// 写入审计数据
tableEnv.executeSql("INSERT INTO audit_table SELECT nextPointer, processTime, status FROM audit_view");

内容的提问来源于stack exchange,提问作者user2386966

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 18:57:43