This commit is contained in:
Rossen Stoyanchev 2016-01-07 15:26:11 -05:00
parent c3a8bf4d17
commit a712f43654
2 changed files with 4 additions and 4 deletions

View File

@ -23,7 +23,7 @@ import java.util.concurrent.BlockingQueue;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscription;
import reactor.rx.Streams;
import reactor.rx.Stream;
/**
* {@code InputStream} implementation based on a byte array {@link Publisher}.
@ -60,7 +60,7 @@ public class ByteBufferPublisherInputStream extends InputStream {
public ByteBufferPublisherInputStream(Publisher<ByteBuffer> publisher, int requestSize) {
Assert.notNull(publisher, "'publisher' must not be null");
this.queue = Streams.from(publisher).toBlockingQueue(requestSize);
this.queue = Stream.from(publisher).toBlockingQueue(requestSize);
}

View File

@ -22,9 +22,9 @@ import javax.xml.bind.Unmarshaller;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import reactor.Flux;
import reactor.Mono;
import reactor.io.buffer.Buffer;
import reactor.rx.Streams;
import org.springframework.http.MediaType;
import org.springframework.util.BufferOutputStream;
@ -73,7 +73,7 @@ public class XmlHandler implements HttpHandler {
bos.close();
buffer.flip();
return response.setBody(Streams.just(buffer.byteBuffer()));
return response.setBody(Flux.just(buffer.byteBuffer()));
}
catch (Exception ex) {
logger.error(ex, ex);