能否实现以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
相关产品推荐
相关产品推荐

