Reactor背压处理:压力缓解前暂停数据库读取
解决方案:基于Pull模式的背压处理+退避重试
你的核心问题在于使用了Push模式的Flux.create,它会无视下游需求持续推送数据库数据,导致背压时要么内存溢出要么直接终止流。正确的做法是改用Pull模式的Flux.generate来响应下游背压,结合重试退避处理Kinesis发送失败。
关键改进点
- 用Flux.generate实现响应式读取:下游请求多少数据,才从数据库读取多少,天然适配背压,不会出现内存积压
- 数据库资源安全管理:在流终止时自动关闭ResultSet、Statement、Connection
- 退避重试机制:针对Kinesis发送失败实现指数退避重试,避免重试风暴
完整代码示例
import reactor.core.publisher.Flux; import reactor.core.scheduler.Schedulers; import reactor.util.retry.Retry; import software.amazon.awssdk.services.kinesis.model.PutRecordRequest; import java.nio.ByteBuffer; import java.sql.*; import java.util.ArrayList; import java.util.HashMap; import java.util.Map; import java.time.Duration; import java.io.IOException; // 封装数据库资源,用于在Flux.generate中管理状态 class ResultSetState implements AutoCloseable { private final Connection conn; private final Statement stmt; private final ResultSet rs; public ResultSetState(Connection conn, Statement stmt, ResultSet rs) { this.conn = conn; this.stmt = stmt; this.rs = rs; } public boolean hasNext() throws SQLException { return rs.next(); } public Map<String, List<Object>> getCurrentValue() throws SQLException { Map<String, List<Object>> value = new HashMap<>(); int columnCount = rs.getMetaData().getColumnCount(); for (int i = 1; i <= columnCount; i++) { String columnName = rs.getMetaData().getColumnName(i); value.computeIfAbsent(columnName, k -> new ArrayList<>()).add(rs.getObject(i)); } return value; } @Override public void close() throws Exception { // 按逆序关闭资源 rs.close(); stmt.close(); conn.close(); } } // 业务逻辑实现 public class DbToKinesisService { // 假设已初始化Kinesis客户端 private final software.amazon.awssdk.services.kinesis.KinesisClient kinesisClient; public void startTransfer() { Flux.generate( // 初始化:建立数据库连接并执行查询 () -> { Connection conn = DriverManager.getConnection("jdbc:your-db-url", "user", "password"); Statement stmt = conn.createStatement(); ResultSet rs = stmt.executeQuery("SELECT * FROM your_target_table"); return new ResultSetState(conn, stmt, rs); }, // 生成元素:仅当下游请求时,读取下一条记录 (state, sink) -> { try { if (state.hasNext()) { sink.next(state.getCurrentValue()); } else { sink.complete(); // 结果集遍历完成,结束流 } } catch (SQLException e) { sink.error(e); // 数据库查询错误,终止流 } }, // 清理:流终止时关闭所有数据库资源 ResultSetState::close ) .subscribeOn(Schedulers.boundedElastic()) // 数据库操作放在单独线程池 .handle((message, sink) -> { try { // 发送数据到Kinesis kinesisClient.putRecord(PutRecordRequest.builder() .streamName("your-kinesis-stream") .data(ByteBuffer.wrap(serialize(message))) // 自行实现序列化逻辑 .partitionKey("default-partition") .build()); sink.next(message); // 发送成功,传递元素 } catch (Exception e) { sink.error(e); // 发送失败,传递错误给重试机制 } }) .retryWhen(Retry.backoff(5, Duration.ofSeconds(1)) // 最多重试5次,初始退避1秒 .jitter(0.5) // 添加50%随机抖动,避免重试风暴 .filter(e -> { // 仅重试Kinesis相关的可恢复错误(如限流、超时) return e instanceof software.amazon.awssdk.core.exception.SdkException || e instanceof IOException; }) .onRetryExhaustedThrow((spec, signal) -> signal.failure())) // 重试耗尽后抛出原始错误 .subscribe( success -> {}, // 成功处理的回调(可按需记录日志) error -> System.err.println("最终处理失败:" + error.getMessage()) ); } // 示例序列化方法,可替换为JSON/Protobuf等 private byte[] serialize(Map<String, List<Object>> data) throws IOException { // 自行实现序列化逻辑 return new byte[0]; } }
为什么这个方案有效
- 背压自动处理:Flux.generate是Pull模式,下游(Kinesis发送逻辑)请求多少数据,上游才从数据库读取多少。当Kinesis积压时,下游需求减少,数据库读取会自动暂停,直到压力缓解。
- 错误重试可控:通过
retryWhen结合Retry.backoff实现指数退避,既保证了错误恢复能力,又避免了频繁重试对Kinesis的冲击。 - 资源安全:通过
AutoCloseable和generate的清理回调,确保数据库资源在流终止(正常结束或异常终止)时被正确关闭,避免资源泄漏。
对原问题的解释
- 你之前用
Flux.create的同步循环是Push模式,不管下游需求持续推送数据,OverflowStrategy.ERROR抛出的异常是由Reactor框架在内部线程抛出的,无法通过try/catch捕获,直接导致整个流终止。 onErrorContinue未生效是因为它仅处理下游操作符(如map、filter)的错误,而Flux.create中的错误属于源错误,不在其处理范围内。改用Flux.generate后,错误会被正确传递到重试机制中。
内容的提问来源于stack exchange,提问作者Ben Shoemaker
相关产品推荐
相关产品推荐

