From 72b66c9715035d8263a55ca726d90426afca1fff Mon Sep 17 00:00:00 2001 From: Arjen Poutsma Date: Tue, 26 Jan 2016 12:33:03 +0100 Subject: [PATCH] Introduction of PooledDataBuffer This commit introduces a pooled data buffer as a subtype of DataBuffer, as well as various utility methods related to reference counting. Additionally, Crelease calls have been introduced throughout the codebase to properly dispose of pooled databuffers. --- .../codec/support/JacksonJsonDecoder.java | 7 +- .../codec/support/JacksonJsonEncoder.java | 19 ++--- .../core/codec/support/JsonObjectDecoder.java | 3 + .../core/codec/support/XmlEventDecoder.java | 1 + .../core/io/buffer/NettyDataBuffer.java | 13 +++- .../core/io/buffer/PooledDataBuffer.java | 40 ++++++++++ .../io/buffer/support/DataBufferUtils.java | 15 +++- .../reactive/ServletHttpHandlerAdapter.java | 3 +- .../core/io/buffer/DataBufferTests.java | 26 ++++++- .../core/io/buffer/PooledDataBufferTests.java | 73 +++++++++++++++++++ 10 files changed, 183 insertions(+), 17 deletions(-) create mode 100644 spring-web-reactive/src/main/java/org/springframework/core/io/buffer/PooledDataBuffer.java create mode 100644 spring-web-reactive/src/test/java/org/springframework/core/io/buffer/PooledDataBufferTests.java diff --git a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonDecoder.java b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonDecoder.java index 05f6dbf2c5..3d6f406f7e 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonDecoder.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonDecoder.java @@ -28,6 +28,7 @@ import org.springframework.core.ResolvableType; import org.springframework.core.codec.CodecException; import org.springframework.core.codec.Decoder; import org.springframework.core.io.buffer.DataBuffer; +import org.springframework.core.io.buffer.support.DataBufferUtils; import org.springframework.util.MimeType; @@ -70,9 +71,11 @@ public class JacksonJsonDecoder extends AbstractDecoder { stream = this.preProcessor.decode(inputStream, type, mimeType, hints); } - return stream.map(content -> { + return stream.map(dataBuffer -> { try { - return reader.readValue(content.asInputStream()); + Object value = reader.readValue(dataBuffer.asInputStream()); + DataBufferUtils.release(dataBuffer); + return value; } catch (IOException e) { throw new CodecException("Error while reading the data", e); diff --git a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonEncoder.java b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonEncoder.java index 01b2bcbc41..fc08cfd878 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonEncoder.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JacksonJsonEncoder.java @@ -18,6 +18,7 @@ package org.springframework.core.codec.support; import java.io.IOException; import java.io.OutputStream; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import com.fasterxml.jackson.databind.ObjectMapper; @@ -43,6 +44,12 @@ public class JacksonJsonEncoder extends AbstractEncoder { private final ObjectMapper mapper; + private static final ByteBuffer START_ARRAY_BUFFER = ByteBuffer.wrap(new byte[]{'['}); + + private static final ByteBuffer SEPARATOR_BUFFER = ByteBuffer.wrap(new byte[]{','}); + + private static final ByteBuffer END_ARRAY_BUFFER = ByteBuffer.wrap(new byte[]{']'}); + public JacksonJsonEncoder() { this(new ObjectMapper()); } @@ -65,10 +72,10 @@ public class JacksonJsonEncoder extends AbstractEncoder { } else { // array - Mono startArray = Mono.just(charBuffer('[', allocator)); + Mono startArray = Mono.just(allocator.wrap(START_ARRAY_BUFFER)); Flux arraySeparators = - Flux.create(sub -> sub.onNext(charBuffer(',', allocator))); - Mono endArray = Mono.just(charBuffer(']', allocator)); + Flux.create(sub -> sub.onNext(allocator.wrap(SEPARATOR_BUFFER))); + Mono endArray = Mono.just(allocator.wrap(END_ARRAY_BUFFER)); Flux serializedObjects = Flux.from(inputStream).map(value -> serialize(value, allocator)); @@ -94,11 +101,5 @@ public class JacksonJsonEncoder extends AbstractEncoder { return buffer; } - private DataBuffer charBuffer(char ch, DataBufferAllocator allocator) { - DataBuffer buffer = allocator.allocateBuffer(1); - buffer.write((byte) ch); - return buffer; - } - } diff --git a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JsonObjectDecoder.java b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JsonObjectDecoder.java index 42d204916a..1c37fe14f3 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JsonObjectDecoder.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/JsonObjectDecoder.java @@ -30,6 +30,7 @@ import reactor.core.publisher.Flux; import org.springframework.core.ResolvableType; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferAllocator; +import org.springframework.core.io.buffer.support.DataBufferUtils; import org.springframework.util.MimeType; /** @@ -110,11 +111,13 @@ public class JsonObjectDecoder extends AbstractDecoder { List chunks = new ArrayList<>(); if (this.input == null) { this.input = Unpooled.copiedBuffer(b.asByteBuffer()); + DataBufferUtils.release(b); this.writerIndex = this.input.writerIndex(); } else { this.input = Unpooled.copiedBuffer(this.input, Unpooled.copiedBuffer(b.asByteBuffer())); + DataBufferUtils.release(b); this.writerIndex = this.input.writerIndex(); } if (this.state == ST_CORRUPTED) { diff --git a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/XmlEventDecoder.java b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/XmlEventDecoder.java index 1a2a77ca60..36a1badb28 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/codec/support/XmlEventDecoder.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/codec/support/XmlEventDecoder.java @@ -134,6 +134,7 @@ public class XmlEventDecoder extends AbstractDecoder { } } } + DataBufferUtils.release(dataBuffer); return Flux.fromIterable(events); } catch (XMLStreamException ex) { diff --git a/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/NettyDataBuffer.java b/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/NettyDataBuffer.java index 8beb2f41b7..7078142d6e 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/NettyDataBuffer.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/NettyDataBuffer.java @@ -36,7 +36,7 @@ import org.springframework.util.ObjectUtils; * * @author Arjen Poutsma */ -public class NettyDataBuffer implements DataBuffer { +public class NettyDataBuffer implements PooledDataBuffer { private final NettyDataBufferAllocator allocator; @@ -181,6 +181,17 @@ public class NettyDataBuffer implements DataBuffer { return new ByteBufOutputStream(this.byteBuf); } + @Override + public PooledDataBuffer retain() { + this.byteBuf.retain(); + return this; + } + + @Override + public boolean release() { + return this.byteBuf.release(); + } + @Override public int hashCode() { return this.byteBuf.hashCode(); diff --git a/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/PooledDataBuffer.java b/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/PooledDataBuffer.java new file mode 100644 index 0000000000..ae2b2f7ca0 --- /dev/null +++ b/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/PooledDataBuffer.java @@ -0,0 +1,40 @@ +/* + * Copyright 2002-2016 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.core.io.buffer; + +/** + * Extension of {@link DataBuffer} that allows for buffer that share a memory pool. + * Introduces methods for reference counting. + * + * @author Arjen Poutsma + */ +public interface PooledDataBuffer extends DataBuffer { + + /** + * Increases the reference count for this buffer by one. + * @return this buffer + */ + PooledDataBuffer retain(); + + /** + * Decreases the reference count for this buffer by one, and releases it once the + * count reaches zero. + * @return {@code true} if the buffer was released; {@code false} otherwise. + */ + boolean release(); + +} diff --git a/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/support/DataBufferUtils.java b/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/support/DataBufferUtils.java index 3d41b007bf..2bb2bcdf0c 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/support/DataBufferUtils.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/io/buffer/support/DataBufferUtils.java @@ -35,6 +35,7 @@ import reactor.core.subscriber.SubscriberWithContext; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferAllocator; +import org.springframework.core.io.buffer.PooledDataBuffer; import org.springframework.util.Assert; import org.springframework.util.CollectionUtils2; @@ -170,6 +171,18 @@ public abstract class DataBufferUtils { }); } + /** + * Releases the given data buffer, if it is a {@link PooledDataBuffer}. + * @param dataBuffer the data buffer to release + * @return {@code true} if the buffer was released; {@code false} otherwise. + */ + public static boolean release(DataBuffer dataBuffer) { + if (dataBuffer instanceof PooledDataBuffer) { + return ((PooledDataBuffer) dataBuffer).release(); + } + return false; + } + private static class ReadableByteChannelConsumer implements Consumer> { @@ -199,7 +212,7 @@ public abstract class DataBufferUtils { } finally { if (release) { - // TODO: release buffer when we have PooledDataBuffer + release(dataBuffer); } } } diff --git a/spring-web-reactive/src/main/java/org/springframework/http/server/reactive/ServletHttpHandlerAdapter.java b/spring-web-reactive/src/main/java/org/springframework/http/server/reactive/ServletHttpHandlerAdapter.java index 329949bee5..f65ec1a16a 100644 --- a/spring-web-reactive/src/main/java/org/springframework/http/server/reactive/ServletHttpHandlerAdapter.java +++ b/spring-web-reactive/src/main/java/org/springframework/http/server/reactive/ServletHttpHandlerAdapter.java @@ -39,6 +39,7 @@ import reactor.core.util.BackpressureUtils; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferAllocator; import org.springframework.core.io.buffer.DefaultDataBufferAllocator; +import org.springframework.core.io.buffer.support.DataBufferUtils; import org.springframework.http.HttpStatus; import org.springframework.util.Assert; @@ -361,7 +362,7 @@ public class ServletHttpHandlerAdapter extends HttpServlet { } private void releaseBuffer() { - // TODO: call PooledDataBuffer.release() when we it is introduced + DataBufferUtils.release(dataBuffer); dataBuffer = null; } diff --git a/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/DataBufferTests.java b/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/DataBufferTests.java index 8169f5e142..687e485e1f 100644 --- a/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/DataBufferTests.java +++ b/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/DataBufferTests.java @@ -28,6 +28,8 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.Parameterized; +import org.springframework.core.io.buffer.support.DataBufferUtils; + import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; @@ -56,6 +58,10 @@ public class DataBufferTests { return allocator.allocateBuffer(capacity); } + private void release(DataBuffer... buffers) { + Arrays.stream(buffers).forEach(DataBufferUtils::release); + } + @Test public void writeAndRead() { @@ -72,6 +78,8 @@ public class DataBufferTests { buffer.read(result); assertArrayEquals(new byte[]{'b', 'c', 'd', 'e'}, result); + + release(buffer); } @Test @@ -103,6 +111,8 @@ public class DataBufferTests { len = inputStream.read(bytes); assertEquals(1, len); assertArrayEquals(new byte[]{'e', (byte) 0}, bytes); + + release(buffer); } @Test @@ -118,6 +128,8 @@ public class DataBufferTests { byte[] bytes = new byte[5]; buffer.read(bytes); assertArrayEquals(new byte[]{'a', 'b', 'c', 'd', 'e'}, bytes); + + release(buffer); } @Test @@ -135,6 +147,8 @@ public class DataBufferTests { result = new byte[2]; buffer.read(result); assertArrayEquals(new byte[]{'c', 'd'}, result); + + release(buffer); } @Test @@ -156,6 +170,12 @@ public class DataBufferTests { buffer1.read(result); assertArrayEquals(new byte[]{'a', 'b', 'c', 'd'}, result); + + release(buffer1); + } + + private ByteBuffer createByteBuffer(int capacity) { + return ByteBuffer.allocate(capacity); } @Test @@ -175,10 +195,8 @@ public class DataBufferTests { buffer1.read(result); assertArrayEquals(new byte[]{'a', 'b', 'c', 'd'}, result); - } - private ByteBuffer createByteBuffer(int capacity) { - return ByteBuffer.allocate(capacity); + release(buffer1); } @Test @@ -195,6 +213,8 @@ public class DataBufferTests { buffer.read(resultBytes); assertArrayEquals(new byte[]{'b', 'c'}, resultBytes); + release(buffer); + } diff --git a/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/PooledDataBufferTests.java b/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/PooledDataBufferTests.java new file mode 100644 index 0000000000..c0f4ae591f --- /dev/null +++ b/spring-web-reactive/src/test/java/org/springframework/core/io/buffer/PooledDataBufferTests.java @@ -0,0 +1,73 @@ +/* + * Copyright 2002-2016 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.core.io.buffer; + +import io.netty.buffer.PooledByteBufAllocator; +import io.netty.buffer.UnpooledByteBufAllocator; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.Parameterized; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +/** + * @author Arjen Poutsma + */ +@RunWith(Parameterized.class) +public class PooledDataBufferTests { + + @Parameterized.Parameter + public DataBufferAllocator allocator; + + @Parameterized.Parameters(name = "{0}") + public static Object[][] buffers() { + + return new Object[][]{ + {new NettyDataBufferAllocator(new UnpooledByteBufAllocator(true))}, + {new NettyDataBufferAllocator(new UnpooledByteBufAllocator(false))}, + {new NettyDataBufferAllocator(new PooledByteBufAllocator(true))}, + {new NettyDataBufferAllocator(new PooledByteBufAllocator(false))}}; + } + + private PooledDataBuffer createDataBuffer(int capacity) { + return (PooledDataBuffer) allocator.allocateBuffer(capacity); + } + + @Test + public void retainAndRelease() { + PooledDataBuffer buffer = createDataBuffer(1); + buffer.write((byte) 'a'); + + buffer.retain(); + boolean result = buffer.release(); + assertFalse(result); + result = buffer.release(); + assertTrue(result); + } + + @Test(expected = IllegalStateException.class) + public void tooManyReleases() { + PooledDataBuffer buffer = createDataBuffer(1); + buffer.write((byte) 'a'); + + buffer.release(); + buffer.release(); + } + + +} \ No newline at end of file