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

Reactor背压处理:压力缓解前暂停数据库读取

解决方案:基于Pull模式的背压处理+退避重试

你的核心问题在于使用了Push模式的Flux.create,它会无视下游需求持续推送数据库数据,导致背压时要么内存溢出要么直接终止流。正确的做法是改用Pull模式的Flux.generate来响应下游背压,结合重试退避处理Kinesis发送失败。

关键改进点

  1. 用Flux.generate实现响应式读取:下游请求多少数据,才从数据库读取多少,天然适配背压,不会出现内存积压
  2. 数据库资源安全管理:在流终止时自动关闭ResultSet、Statement、Connection
  3. 退避重试机制:针对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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 10:22:01