Skip to content

Commit d2f663d

Browse files
GH-1194: Preserve empty list offset buffers
1 parent 196b58a commit d2f663d

5 files changed

Lines changed: 263 additions & 43 deletions

File tree

vector/src/main/java/org/apache/arrow/vector/complex/LargeListVector.java

Lines changed: 46 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -305,18 +305,41 @@ public void exportCDataBuffers(List<ArrowBuf> buffers, ArrowBuf buffersPtr, long
305305

306306
/** Set the reader and writer indexes for the inner buffers. */
307307
private void setReaderAndWriterIndex() {
308+
final long requiredOffsetBufferCapacity = (long) (valueCount + 1) * OFFSET_WIDTH;
308309
validityBuffer.readerIndex(0);
309310
offsetBuffer.readerIndex(0);
310311
if (valueCount == 0) {
311312
validityBuffer.writerIndex(0);
313+
ensureEmptyOffsetBufferCapacity(requiredOffsetBufferCapacity);
312314
} else {
313315
validityBuffer.writerIndex(BitVectorHelper.getValidityBufferSizeFromCount(valueCount));
314316
}
315317
// IPC serializer will determine readable bytes based on `readerIndex` and `writerIndex`.
316318
// Both are set to 0 means 0 bytes are written to the IPC stream which will crash IPC readers
317319
// in other libraries. According to Arrow spec, we should still output the offset buffer which
318320
// is [0].
319-
offsetBuffer.writerIndex((long) (valueCount + 1) * OFFSET_WIDTH);
321+
offsetBuffer.writerIndex(requiredOffsetBufferCapacity);
322+
}
323+
324+
private void ensureEmptyOffsetBufferCapacity(long requiredCapacity) {
325+
if (offsetBuffer.capacity() >= requiredCapacity) {
326+
return;
327+
}
328+
long previousOffsetAllocationSizeInBytes = offsetAllocationSizeInBytes;
329+
ArrowBuf oldOffsetBuffer = offsetBuffer;
330+
offsetBuffer = allocateOffsetBuffer(requiredCapacity);
331+
final long bytesToCopy = Math.min(oldOffsetBuffer.capacity(), requiredCapacity);
332+
offsetBuffer.setBytes(0, oldOffsetBuffer, 0, bytesToCopy);
333+
334+
final int copiedOffsets = (int) (bytesToCopy / OFFSET_WIDTH);
335+
final int requiredOffsets = (int) (requiredCapacity / OFFSET_WIDTH);
336+
final long lastCopiedOffset =
337+
copiedOffsets == 0 ? 0 : offsetBuffer.getLong((long) (copiedOffsets - 1) * OFFSET_WIDTH);
338+
for (int i = copiedOffsets; i < requiredOffsets; i++) {
339+
offsetBuffer.setLong((long) i * OFFSET_WIDTH, lastCopiedOffset);
340+
}
341+
offsetAllocationSizeInBytes = previousOffsetAllocationSizeInBytes;
342+
oldOffsetBuffer.getReferenceManager().release();
320343
}
321344

322345
/**
@@ -674,24 +697,30 @@ public void splitAndTransfer(int startIndex, int length) {
674697
startIndex,
675698
length,
676699
valueCount);
677-
final long startPoint = offsetBuffer.getLong((long) startIndex * OFFSET_WIDTH);
678-
final long sliceLength =
679-
offsetBuffer.getLong((long) (startIndex + length) * OFFSET_WIDTH) - startPoint;
680700
to.clear();
681-
to.offsetBuffer = to.allocateOffsetBuffer((length + 1) * OFFSET_WIDTH);
682-
/* splitAndTransfer offset buffer */
683-
for (int i = 0; i < length + 1; i++) {
684-
final long relativeOffset =
685-
offsetBuffer.getLong((long) (startIndex + i) * OFFSET_WIDTH) - startPoint;
686-
to.offsetBuffer.setLong((long) i * OFFSET_WIDTH, relativeOffset);
701+
if (length > 0) {
702+
final long startPoint = offsetBuffer.getLong((long) startIndex * OFFSET_WIDTH);
703+
final long sliceLength =
704+
offsetBuffer.getLong((long) (startIndex + length) * OFFSET_WIDTH) - startPoint;
705+
to.offsetBuffer = to.allocateOffsetBuffer((length + 1) * OFFSET_WIDTH);
706+
/* splitAndTransfer offset buffer */
707+
for (int i = 0; i < length + 1; i++) {
708+
final long relativeOffset =
709+
offsetBuffer.getLong((long) (startIndex + i) * OFFSET_WIDTH) - startPoint;
710+
to.offsetBuffer.setLong((long) i * OFFSET_WIDTH, relativeOffset);
711+
}
712+
/* splitAndTransfer validity buffer */
713+
splitAndTransferValidityBuffer(startIndex, length, to);
714+
/* splitAndTransfer data buffer */
715+
dataTransferPair.splitAndTransfer(
716+
checkedCastToInt(startPoint), checkedCastToInt(sliceLength));
717+
to.lastSet = length - 1;
718+
to.setValueCount(length);
719+
} else {
720+
to.ensureEmptyOffsetBufferCapacity(OFFSET_WIDTH);
721+
dataTransferPair.splitAndTransfer(0, 0);
722+
to.setValueCount(0);
687723
}
688-
/* splitAndTransfer validity buffer */
689-
splitAndTransferValidityBuffer(startIndex, length, to);
690-
/* splitAndTransfer data buffer */
691-
dataTransferPair.splitAndTransfer(
692-
checkedCastToInt(startPoint), checkedCastToInt(sliceLength));
693-
to.lastSet = length - 1;
694-
to.setValueCount(length);
695724
}
696725

697726
@Override

vector/src/main/java/org/apache/arrow/vector/complex/ListVector.java

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -263,18 +263,41 @@ public void exportCDataBuffers(List<ArrowBuf> buffers, ArrowBuf buffersPtr, long
263263

264264
/** Set the reader and writer indexes for the inner buffers. */
265265
private void setReaderAndWriterIndex() {
266+
final long requiredOffsetBufferCapacity = (long) (valueCount + 1) * OFFSET_WIDTH;
266267
validityBuffer.readerIndex(0);
267268
offsetBuffer.readerIndex(0);
268269
if (valueCount == 0) {
269270
validityBuffer.writerIndex(0);
271+
ensureEmptyOffsetBufferCapacity(requiredOffsetBufferCapacity);
270272
} else {
271273
validityBuffer.writerIndex(BitVectorHelper.getValidityBufferSizeFromCount(valueCount));
272274
}
273275
// IPC serializer will determine readable bytes based on `readerIndex` and `writerIndex`.
274276
// Both are set to 0 means 0 bytes are written to the IPC stream which will crash IPC readers
275277
// in other libraries. According to Arrow spec, we should still output the offset buffer which
276278
// is [0].
277-
offsetBuffer.writerIndex((long) (valueCount + 1) * OFFSET_WIDTH);
279+
offsetBuffer.writerIndex(requiredOffsetBufferCapacity);
280+
}
281+
282+
private void ensureEmptyOffsetBufferCapacity(long requiredCapacity) {
283+
if (offsetBuffer.capacity() >= requiredCapacity) {
284+
return;
285+
}
286+
long previousOffsetAllocationSizeInBytes = offsetAllocationSizeInBytes;
287+
ArrowBuf oldOffsetBuffer = offsetBuffer;
288+
offsetBuffer = allocateOffsetBuffer(requiredCapacity);
289+
final long bytesToCopy = Math.min(oldOffsetBuffer.capacity(), requiredCapacity);
290+
offsetBuffer.setBytes(0, oldOffsetBuffer, 0, bytesToCopy);
291+
292+
final int copiedOffsets = (int) (bytesToCopy / OFFSET_WIDTH);
293+
final int requiredOffsets = (int) (requiredCapacity / OFFSET_WIDTH);
294+
final int lastCopiedOffset =
295+
copiedOffsets == 0 ? 0 : offsetBuffer.getInt((copiedOffsets - 1) * OFFSET_WIDTH);
296+
for (int i = copiedOffsets; i < requiredOffsets; i++) {
297+
offsetBuffer.setInt(i * OFFSET_WIDTH, lastCopiedOffset);
298+
}
299+
offsetAllocationSizeInBytes = previousOffsetAllocationSizeInBytes;
300+
oldOffsetBuffer.getReferenceManager().release();
278301
}
279302

280303
/**
@@ -572,6 +595,10 @@ public void splitAndTransfer(int startIndex, int length) {
572595
dataTransferPair.splitAndTransfer(startPoint, sliceLength);
573596
to.lastSet = length - 1;
574597
to.setValueCount(length);
598+
} else {
599+
to.ensureEmptyOffsetBufferCapacity(OFFSET_WIDTH);
600+
dataTransferPair.splitAndTransfer(0, 0);
601+
to.setValueCount(0);
575602
}
576603
}
577604

vector/src/test/java/org/apache/arrow/vector/TestLargeListVector.java

Lines changed: 75 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1102,24 +1102,89 @@ public void testCopyValueSafeForExtensionType() throws Exception {
11021102

11031103
@Test
11041104
public void testEmptyLargeListOffsetBuffer() {
1105-
// Test that LargeListVector has correct readableBytes after allocation.
1106-
// According to Arrow spec, offset buffer must have N+1 entries.
1107-
// Even when N=0, it should contain [0].
11081105
try (LargeListVector list = LargeListVector.empty("list", allocator)) {
11091106
list.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
11101107
list.allocateNew();
11111108
list.setValueCount(0);
11121109

1113-
List<ArrowBuf> buffers = list.getFieldBuffers();
1114-
assertTrue(
1115-
buffers.get(1).readableBytes() >= LargeListVector.OFFSET_WIDTH,
1116-
"Offset buffer should have at least "
1117-
+ LargeListVector.OFFSET_WIDTH
1118-
+ " bytes for offset[0]");
1119-
assertEquals(0L, list.getOffsetBuffer().getLong(0));
1110+
assertEmptyLargeListOffsetBuffer(list);
11201111
}
11211112
}
11221113

1114+
@Test
1115+
public void testUnallocatedEmptyLargeListOffsetBuffer() {
1116+
try (LargeListVector list = LargeListVector.empty("list", allocator)) {
1117+
list.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1118+
list.setValueCount(0);
1119+
1120+
assertEmptyLargeListOffsetBuffer(list);
1121+
}
1122+
}
1123+
1124+
@Test
1125+
public void testSplitAndTransferEmptyLargeListOffsetBuffer() {
1126+
try (LargeListVector source = LargeListVector.empty("source", allocator);
1127+
LargeListVector target = LargeListVector.empty("target", allocator)) {
1128+
source.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1129+
target.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1130+
source.allocateNew();
1131+
source.setValueCount(0);
1132+
1133+
TransferPair transferPair = source.makeTransferPair(target);
1134+
transferPair.splitAndTransfer(0, 0);
1135+
1136+
assertEmptyLargeListOffsetBuffer(target);
1137+
}
1138+
}
1139+
1140+
@Test
1141+
public void testSplitAndTransferEmptyLargeListAllocatesOffsetBuffer() {
1142+
try (LargeListVector fromVector = LargeListVector.empty("fromVector", allocator);
1143+
LargeListVector toVector = LargeListVector.empty("toVector", allocator)) {
1144+
fromVector.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1145+
fromVector.allocateNew();
1146+
fromVector.setValueCount(0);
1147+
1148+
TransferPair transferPair = fromVector.makeTransferPair(toVector);
1149+
transferPair.splitAndTransfer(0, 0);
1150+
1151+
assertAllocatedEmptyLargeListOffsetBuffer(toVector);
1152+
}
1153+
}
1154+
1155+
@Test
1156+
public void testSplitAndTransferEmptyNestedLargeListAllocatesOffsetBuffers() {
1157+
try (LargeListVector fromVector = LargeListVector.empty("fromVector", allocator);
1158+
LargeListVector toVector = LargeListVector.empty("toVector", allocator)) {
1159+
fromVector.addOrGetVector(FieldType.nullable(MinorType.LARGELIST.getType()));
1160+
LargeListVector childVector = (LargeListVector) fromVector.getDataVector();
1161+
childVector.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1162+
fromVector.allocateNew();
1163+
fromVector.setValueCount(0);
1164+
1165+
TransferPair transferPair = fromVector.makeTransferPair(toVector);
1166+
transferPair.splitAndTransfer(0, 0);
1167+
1168+
assertAllocatedEmptyLargeListOffsetBuffer(toVector);
1169+
assertAllocatedEmptyLargeListOffsetBuffer((LargeListVector) toVector.getDataVector());
1170+
}
1171+
}
1172+
1173+
private ArrowBuf assertEmptyLargeListOffsetBuffer(LargeListVector list) {
1174+
List<ArrowBuf> buffers = list.getFieldBuffers();
1175+
ArrowBuf offsetBuffer = buffers.get(1);
1176+
assertEquals(LargeListVector.OFFSET_WIDTH, offsetBuffer.readableBytes());
1177+
assertTrue(offsetBuffer.capacity() >= LargeListVector.OFFSET_WIDTH);
1178+
assertEquals(0L, offsetBuffer.getLong(0));
1179+
return offsetBuffer;
1180+
}
1181+
1182+
private void assertAllocatedEmptyLargeListOffsetBuffer(LargeListVector list) {
1183+
ArrowBuf offsetBuffer = list.getOffsetBuffer();
1184+
assertTrue(offsetBuffer.capacity() >= LargeListVector.OFFSET_WIDTH);
1185+
assertEquals(0L, offsetBuffer.getLong(0));
1186+
}
1187+
11231188
private void writeIntValues(UnionLargeListWriter writer, int[] values) {
11241189
writer.startList();
11251190
for (int v : values) {

vector/src/test/java/org/apache/arrow/vector/TestListVector.java

Lines changed: 75 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1381,24 +1381,89 @@ public void testCopyValueSafeForExtensionType() throws Exception {
13811381

13821382
@Test
13831383
public void testEmptyListOffsetBuffer() {
1384-
// Test that ListVector has correct readableBytes after allocation.
1385-
// According to Arrow spec, offset buffer must have N+1 entries.
1386-
// Even when N=0, it should contain [0].
13871384
try (ListVector list = ListVector.empty("list", allocator)) {
13881385
list.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
13891386
list.allocateNew();
13901387
list.setValueCount(0);
13911388

1392-
List<ArrowBuf> buffers = list.getFieldBuffers();
1393-
assertTrue(
1394-
buffers.get(1).readableBytes() >= BaseRepeatedValueVector.OFFSET_WIDTH,
1395-
"Offset buffer should have at least "
1396-
+ BaseRepeatedValueVector.OFFSET_WIDTH
1397-
+ " bytes for offset[0]");
1398-
assertEquals(0, list.getOffsetBuffer().getInt(0));
1389+
assertEmptyListOffsetBuffer(list);
13991390
}
14001391
}
14011392

1393+
@Test
1394+
public void testUnallocatedEmptyListOffsetBuffer() {
1395+
try (ListVector list = ListVector.empty("list", allocator)) {
1396+
list.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1397+
list.setValueCount(0);
1398+
1399+
assertEmptyListOffsetBuffer(list);
1400+
}
1401+
}
1402+
1403+
@Test
1404+
public void testSplitAndTransferEmptyListOffsetBuffer() {
1405+
try (ListVector source = ListVector.empty("source", allocator);
1406+
ListVector target = ListVector.empty("target", allocator)) {
1407+
source.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1408+
target.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1409+
source.allocateNew();
1410+
source.setValueCount(0);
1411+
1412+
TransferPair transferPair = source.makeTransferPair(target);
1413+
transferPair.splitAndTransfer(0, 0);
1414+
1415+
assertEmptyListOffsetBuffer(target);
1416+
}
1417+
}
1418+
1419+
@Test
1420+
public void testSplitAndTransferEmptyListAllocatesOffsetBuffer() {
1421+
try (ListVector fromVector = ListVector.empty("fromVector", allocator);
1422+
ListVector toVector = ListVector.empty("toVector", allocator)) {
1423+
fromVector.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1424+
fromVector.allocateNew();
1425+
fromVector.setValueCount(0);
1426+
1427+
TransferPair transferPair = fromVector.makeTransferPair(toVector);
1428+
transferPair.splitAndTransfer(0, 0);
1429+
1430+
assertAllocatedEmptyListOffsetBuffer(toVector);
1431+
}
1432+
}
1433+
1434+
@Test
1435+
public void testSplitAndTransferEmptyNestedListAllocatesOffsetBuffers() {
1436+
try (ListVector fromVector = ListVector.empty("fromVector", allocator);
1437+
ListVector toVector = ListVector.empty("toVector", allocator)) {
1438+
fromVector.addOrGetVector(FieldType.nullable(MinorType.LIST.getType()));
1439+
ListVector childVector = (ListVector) fromVector.getDataVector();
1440+
childVector.addOrGetVector(FieldType.nullable(MinorType.INT.getType()));
1441+
fromVector.allocateNew();
1442+
fromVector.setValueCount(0);
1443+
1444+
TransferPair transferPair = fromVector.makeTransferPair(toVector);
1445+
transferPair.splitAndTransfer(0, 0);
1446+
1447+
assertAllocatedEmptyListOffsetBuffer(toVector);
1448+
assertAllocatedEmptyListOffsetBuffer((ListVector) toVector.getDataVector());
1449+
}
1450+
}
1451+
1452+
private ArrowBuf assertEmptyListOffsetBuffer(ListVector list) {
1453+
List<ArrowBuf> buffers = list.getFieldBuffers();
1454+
ArrowBuf offsetBuffer = buffers.get(1);
1455+
assertEquals(BaseRepeatedValueVector.OFFSET_WIDTH, offsetBuffer.readableBytes());
1456+
assertTrue(offsetBuffer.capacity() >= BaseRepeatedValueVector.OFFSET_WIDTH);
1457+
assertEquals(0, offsetBuffer.getInt(0));
1458+
return offsetBuffer;
1459+
}
1460+
1461+
private void assertAllocatedEmptyListOffsetBuffer(ListVector list) {
1462+
ArrowBuf offsetBuffer = list.getOffsetBuffer();
1463+
assertTrue(offsetBuffer.capacity() >= BaseRepeatedValueVector.OFFSET_WIDTH);
1464+
assertEquals(0, offsetBuffer.getInt(0));
1465+
}
1466+
14021467
private void writeIntValues(UnionListWriter writer, int[] values) {
14031468
writer.startList();
14041469
for (int v : values) {

0 commit comments

Comments
 (0)