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

能否实现以Consumer<OutputStream>为支撑的org.springframework.http.HttpEntity?

问题背景

我的目标是通过Spring WebClient传输大于2GB的Blob,问题在于这些Blob通常并非存储在文件系统中,而是位于网络资源(例如数据库)内。因此大多数情况下,我需要先打开与网络资源的会话,并在数据传输完成后释放该会话。

获取InputStream并执行资源清理的实现

@Override
protected InputStream doCreateInputStream(S store, DataTicket dataTicket, Map<String, ContentAttr<?>> contentAttrs) {
    String storeName = store.getName();
    IDfId storeId = store.getObjectId();
    String query = "select content from " + requireValidStoreName(storeName) + " where ticket=:ticket";

    DisposeChain disposeChain = new DisposeChain(null);
    SystemEventLogger systemEventLogger = getSystemEventLogger();

    Supplier<InputStream> supplier = () -> {
        Connection con = DataSourceUtils.getConnection(dataSource);
        disposeChain.prepend(() -> DataSourceUtils.releaseConnection(con, dataSource));
        PreparedStatement stmt;
        boolean success = false;
        try {
            stmt = con.prepareStatement(query);
            disposeChain.prepend(() -> closeSilently(stmt, systemEventLogger));
            stmt.setLong(1, dataTicket.getValue());
            ResultSet resultSet = stmt.executeQuery();
            disposeChain.prepend(() -> closeSilently(resultSet, systemEventLogger));

            if (!resultSet.next()) {
                throw ticketNotFound(storeId, dataTicket);
            }

            LobHandler lobHandler = new DefaultLobHandler();
            InputStream is = lobHandler.getBlobAsBinaryStream(resultSet, "content");
            disposeChain.prepend(() -> closeSilently(is, systemEventLogger));
            success = true;
            return is;
        } catch (Exception ex) {
            throw unknownStoreError(storeId, ex);
        } finally {
            if (!success) {
                disposeChain.run();
            }
        }
    };

    InputStream result = new LazyInputStream(supplier, disposeChain);
    resourceHousekeeper.register(result, disposeChain);
    return result;
}

基于Consumer的简洁实现

不过,基于Consumer<OutputStream>的实现会简洁得多(另一种可选方案是使用PipedInputStream,但我同样不喜欢这种方式):

@Override
protected Consumer<OutputStream> doCreatePuller(S store, DataTicket dataTicket, Map<String, ContentAttr<?>> contentAttrs) {
    String storeName = store.getName();
    IDfId storeId = store.getObjectId();
    return os -> {
        RowMapper<Void> rowMapper = (rs, rowNum) -> {
            LobHandler lobHandler = new DefaultLobHandler();
            InputStream is = lobHandler.getBlobAsBinaryStream(rs, "content");
            IOUtils.copy(is, os);
            return null;
        };
        try {
            String statement = "select content from " + requireValidStoreName(storeName) + " where ticket=:ticket";
            MapSqlParameterSource parameters = new MapSqlParameterSource();
            parameters.addValue("ticket", dataTicket.getValue());
            return jdbcTemplate.queryForObject(statement, parameters, rowMapper);
        } catch (EmptyResultDataAccessException ex) {
            throw ticketNotFound(storeId, dataTicket);
        }
    };
}

核心问题

能否实现一个类似Apache的org.apache.http.entity.AbstractHttpEntity的org.springframework.http.HttpEntity?示例参考如下:

public class OutputStreamConsumingEntity extends AbstractHttpEntity implements HttpAsyncContentProducer {

    private final Consumer<OutputStream> outputStreamConsumer;

    public OutputStreamConsumingEntity(@NonNull Consumer<OutputStream> outputStreamConsumer) {
        this.outputStreamConsumer = outputStreamConsumer;
    }

...

    @Override
    public void writeTo(OutputStream outStream) throws IOException {
        outputStreamConsumer.accept(outStream);
    }

    @Override
    public boolean isStreaming() {
        return true;
    }
}

内容的提问来源于stack exchange,提问作者Andrey B. Panfilov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 23:39:56