如何在Flink中使用Async I/O调用分页HTTP API并实现审计追踪
基于Flink处理分页API并实现审计追踪的方案
针对你提出的分页API调用和审计追踪需求,以下是具体的Flink实现方案:
一、利用Flink Async I/O调用分页API(批处理/流处理模式)
1. 批处理模式实现
批处理为有界场景,核心是迭代获取分页数据直到API返回的next指针为空,通过带状态的AsyncFunction结合迭代逻辑实现:
- 初始触发:从配置或审计数据库读取初始
next指针(首次执行可设为null),作为作业输入源(如用fromElements生成触发信号)。 - 自定义AsyncFunction逻辑:
- 在
asyncInvoke中,根据当前next指针构造API请求,调用HTTP客户端获取响应。 - 解析响应中的业务数据和新的
next指针。 - 将业务数据输出到下游,同时将新
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逻辑:
- 初始化时从审计数据库读取上次保存的
next指针,存入ValueState。 - 每次触发请求时,从状态中取出
next指针构造API请求。 - 解析响应后更新状态中的
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
相关产品推荐
相关产品推荐

