Reduce memory copies and object creations when acting on ByteBufs

patch by Norman Maurer; reviewed by jasobrown for CASSANDRA-13789
This commit is contained in:
Norman Maurer 2017-08-22 14:15:16 +02:00 committed by Jason Brown
parent a5f70ad4b8
commit 6bea7de504
4 changed files with 41 additions and 20 deletions

View File

@ -1,4 +1,5 @@
4.0 4.0
* Reduce memory copies and object creations when acting on ByteBufs (CASSANDRA-13789)
* simplify mx4j configuration (Cassandra-13578) * simplify mx4j configuration (Cassandra-13578)
* Fix trigger example on 4.0 (CASSANDRA-13796) * Fix trigger example on 4.0 (CASSANDRA-13796)
* force minumum timeout value (CASSANDRA-9375) * force minumum timeout value (CASSANDRA-9375)

View File

@ -429,7 +429,7 @@ public abstract class QueryOptions
if (!flags.isEmpty()) if (!flags.isEmpty())
{ {
int pageSize = flags.contains(Flag.PAGE_SIZE) ? body.readInt() : -1; int pageSize = flags.contains(Flag.PAGE_SIZE) ? body.readInt() : -1;
PagingState pagingState = flags.contains(Flag.PAGING_STATE) ? PagingState.deserialize(CBUtil.readValue(body), version) : null; PagingState pagingState = flags.contains(Flag.PAGING_STATE) ? PagingState.deserialize(CBUtil.readValueNoCopy(body), version) : null;
ConsistencyLevel serialConsistency = flags.contains(Flag.SERIAL_CONSISTENCY) ? CBUtil.readConsistencyLevel(body) : ConsistencyLevel.SERIAL; ConsistencyLevel serialConsistency = flags.contains(Flag.SERIAL_CONSISTENCY) ? CBUtil.readConsistencyLevel(body) : ConsistencyLevel.SERIAL;
long timestamp = Long.MIN_VALUE; long timestamp = Long.MIN_VALUE;
if (flags.contains(Flag.TIMESTAMP)) if (flags.contains(Flag.TIMESTAMP))

View File

@ -310,7 +310,7 @@ public class ResultSet
PagingState state = null; PagingState state = null;
if (flags.contains(Flag.HAS_MORE_PAGES)) if (flags.contains(Flag.HAS_MORE_PAGES))
state = PagingState.deserialize(CBUtil.readValue(body), version); state = PagingState.deserialize(CBUtil.readValueNoCopy(body), version);
if (flags.contains(Flag.NO_METADATA)) if (flags.contains(Flag.NO_METADATA))
return new ResultMetadata(flags, null, columnCount, state); return new ResultMetadata(flags, null, columnCount, state);

View File

@ -58,6 +58,7 @@ public abstract class CBUtil
{ {
public static final boolean USE_HEAP_ALLOCATOR = Boolean.getBoolean(Config.PROPERTY_PREFIX + "netty_use_heap_allocator"); public static final boolean USE_HEAP_ALLOCATOR = Boolean.getBoolean(Config.PROPERTY_PREFIX + "netty_use_heap_allocator");
public static final ByteBufAllocator allocator = USE_HEAP_ALLOCATOR ? new UnpooledByteBufAllocator(false) : new PooledByteBufAllocator(true); public static final ByteBufAllocator allocator = USE_HEAP_ALLOCATOR ? new UnpooledByteBufAllocator(false) : new PooledByteBufAllocator(true);
private static final int UUID_SIZE = 16;
private final static FastThreadLocal<CharsetDecoder> TL_UTF8_DECODER = new FastThreadLocal<CharsetDecoder>() private final static FastThreadLocal<CharsetDecoder> TL_UTF8_DECODER = new FastThreadLocal<CharsetDecoder>()
{ {
@ -139,9 +140,9 @@ public abstract class CBUtil
public static void writeString(String str, ByteBuf cb) public static void writeString(String str, ByteBuf cb)
{ {
int writerIndex = cb.writerIndex(); int writerIndex = cb.writerIndex();
cb.writeShort(0); cb.writerIndex(writerIndex + 2);
int lengthBytes = ByteBufUtil.writeUtf8(cb, str); int written = ByteBufUtil.writeUtf8(cb, str);
cb.setShort(writerIndex, lengthBytes); cb.setShort(writerIndex, written);
} }
public static int sizeOfString(String str) public static int sizeOfString(String str)
@ -164,14 +165,15 @@ public abstract class CBUtil
public static void writeLongString(String str, ByteBuf cb) public static void writeLongString(String str, ByteBuf cb)
{ {
byte[] bytes = str.getBytes(CharsetUtil.UTF_8); int writerIndex = cb.writerIndex();
cb.writeInt(bytes.length); cb.writerIndex(writerIndex + 4);
cb.writeBytes(bytes); int written = ByteBufUtil.writeUtf8(cb, str);
cb.setInt(writerIndex, written);
} }
public static int sizeOfLongString(String str) public static int sizeOfLongString(String str)
{ {
return 4 + str.getBytes(CharsetUtil.UTF_8).length; return 4 + TypeSizes.encodedUTF8Length(str);
} }
public static byte[] readBytes(ByteBuf cb) public static byte[] readBytes(ByteBuf cb)
@ -274,9 +276,9 @@ public abstract class CBUtil
public static UUID readUUID(ByteBuf cb) public static UUID readUUID(ByteBuf cb)
{ {
byte[] bytes = new byte[16]; ByteBuffer buffer = cb.nioBuffer(cb.readerIndex(), UUID_SIZE);
cb.readBytes(bytes); cb.skipBytes(buffer.remaining());
return UUIDGen.getUUID(ByteBuffer.wrap(bytes)); return UUIDGen.getUUID(buffer);
} }
public static void writeUUID(UUID uuid, ByteBuf cb) public static void writeUUID(UUID uuid, ByteBuf cb)
@ -286,7 +288,7 @@ public abstract class CBUtil
public static int sizeOfUUID(UUID uuid) public static int sizeOfUUID(UUID uuid)
{ {
return 16; return UUID_SIZE;
} }
public static List<String> readStringList(ByteBuf cb) public static List<String> readStringList(ByteBuf cb)
@ -386,9 +388,19 @@ public abstract class CBUtil
int length = cb.readInt(); int length = cb.readInt();
if (length < 0) if (length < 0)
return null; return null;
ByteBuf slice = cb.readSlice(length);
return ByteBuffer.wrap(readRawBytes(slice)); return ByteBuffer.wrap(readRawBytes(cb, length));
}
public static ByteBuffer readValueNoCopy(ByteBuf cb)
{
int length = cb.readInt();
if (length < 0)
return null;
ByteBuffer buffer = cb.nioBuffer(cb.readerIndex(), length);
cb.skipBytes(length);
return buffer;
} }
public static ByteBuffer readBoundValue(ByteBuf cb, ProtocolVersion protocolVersion) public static ByteBuffer readBoundValue(ByteBuf cb, ProtocolVersion protocolVersion)
@ -405,9 +417,7 @@ public abstract class CBUtil
else else
throw new ProtocolException("Invalid ByteBuf length " + length); throw new ProtocolException("Invalid ByteBuf length " + length);
} }
ByteBuf slice = cb.readSlice(length); return ByteBuffer.wrap(readRawBytes(cb, length));
return ByteBuffer.wrap(readRawBytes(slice));
} }
public static void writeValue(byte[] bytes, ByteBuf cb) public static void writeValue(byte[] bytes, ByteBuf cb)
@ -434,7 +444,12 @@ public abstract class CBUtil
cb.writeInt(remaining); cb.writeInt(remaining);
if (remaining > 0) if (remaining > 0)
cb.writeBytes(bytes.duplicate()); {
// write and reset original position so we not need to create a duplicate.
int position = bytes.position();
cb.writeBytes(bytes);
bytes.position(position);
}
} }
public static int sizeOfValue(byte[] bytes) public static int sizeOfValue(byte[] bytes)
@ -560,7 +575,12 @@ public abstract class CBUtil
*/ */
public static byte[] readRawBytes(ByteBuf cb) public static byte[] readRawBytes(ByteBuf cb)
{ {
byte[] bytes = new byte[cb.readableBytes()]; return readRawBytes(cb, cb.readableBytes());
}
private static byte[] readRawBytes(ByteBuf cb, int length)
{
byte[] bytes = new byte[length];
cb.readBytes(bytes); cb.readBytes(bytes);
return bytes; return bytes;
} }