Skip to content

Commit 8ce3973

Browse files
jordepicJordan Epstein
andauthored
GH-1179: Correct the size of var-width vector with >0 start offset during vector append (#1180)
## What's Changed Fix VectorAppender data size computation for variable-width vectors with non-zero start offsets When appending a variable width offset vector in DataFusion comet I was receiving exceptions due to allocating too much memory. This is because Comet passes variable width arrays back to Java where the initial offset vector entry is greater than 0. Prior to this change, arrow-java determines how many bytes to copy by just looking at the last offset entry in the buffer, completely disregarding the value of the first. If first = 100 and last = 200, Java will still copy 200 bytes instead of 100. In this change we fix that. Closes #1179 --------- Co-authored-by: Jordan Epstein <jordan.epstein@imc.com>
1 parent 04471f4 commit 8ce3973

2 files changed

Lines changed: 257 additions & 24 deletions

File tree

vector/src/main/java/org/apache/arrow/vector/util/VectorAppender.java

Lines changed: 68 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -125,10 +125,15 @@ public ValueVector visit(BaseVariableWidthVector deltaVector, Void value) {
125125
targetVector
126126
.getOffsetBuffer()
127127
.getInt((long) targetVector.getValueCount() * BaseVariableWidthVector.OFFSET_WIDTH);
128+
// The delta vector's offset buffer need not start at zero (e.g. a vector imported through
129+
// the C data interface from a sliced array), so the amount of data to append is the
130+
// distance between its first and last offsets, not the last offset itself.
131+
int deltaDataStart = deltaVector.getOffsetBuffer().getInt(0);
128132
int deltaDataSize =
129133
deltaVector
130-
.getOffsetBuffer()
131-
.getInt((long) deltaVector.getValueCount() * BaseVariableWidthVector.OFFSET_WIDTH);
134+
.getOffsetBuffer()
135+
.getInt((long) deltaVector.getValueCount() * BaseVariableWidthVector.OFFSET_WIDTH)
136+
- deltaDataStart;
132137
int newValueCapacity = targetDataSize + deltaDataSize;
133138

134139
// make sure there is enough capacity
@@ -149,7 +154,7 @@ public ValueVector visit(BaseVariableWidthVector deltaVector, Void value) {
149154

150155
// append data buffer
151156
MemoryUtil.copyMemory(
152-
deltaVector.getDataBuffer().memoryAddress(),
157+
deltaVector.getDataBuffer().memoryAddress() + deltaDataStart,
153158
targetVector.getDataBuffer().memoryAddress() + targetDataSize,
154159
deltaDataSize);
155160

@@ -160,7 +165,7 @@ public ValueVector visit(BaseVariableWidthVector deltaVector, Void value) {
160165
+ (targetVector.getValueCount() + 1) * BaseVariableWidthVector.OFFSET_WIDTH,
161166
deltaVector.getValueCount() * BaseVariableWidthVector.OFFSET_WIDTH);
162167

163-
// increase each offset from the second buffer
168+
// rebase each appended offset to the target's data, accounting for the delta's start offset
164169
for (int i = 0; i < deltaVector.getValueCount(); i++) {
165170
int oldOffset =
166171
targetVector
@@ -172,7 +177,7 @@ public ValueVector visit(BaseVariableWidthVector deltaVector, Void value) {
172177
.getOffsetBuffer()
173178
.setInt(
174179
(long) (targetVector.getValueCount() + 1 + i) * BaseVariableWidthVector.OFFSET_WIDTH,
175-
oldOffset + targetDataSize);
180+
oldOffset - deltaDataStart + targetDataSize);
176181
}
177182
((BaseVariableWidthVector) targetVector).setLastSet(newValueCount - 1);
178183
targetVector.setValueCount(newValueCount);
@@ -196,11 +201,15 @@ public ValueVector visit(BaseLargeVariableWidthVector deltaVector, Void value) {
196201
.getOffsetBuffer()
197202
.getLong(
198203
(long) targetVector.getValueCount() * BaseLargeVariableWidthVector.OFFSET_WIDTH);
204+
// see the corresponding comment in visit(BaseVariableWidthVector, Void): the delta's
205+
// offset buffer need not start at zero
206+
long deltaDataStart = deltaVector.getOffsetBuffer().getLong(0);
199207
long deltaDataSize =
200208
deltaVector
201-
.getOffsetBuffer()
202-
.getLong(
203-
(long) deltaVector.getValueCount() * BaseLargeVariableWidthVector.OFFSET_WIDTH);
209+
.getOffsetBuffer()
210+
.getLong(
211+
(long) deltaVector.getValueCount() * BaseLargeVariableWidthVector.OFFSET_WIDTH)
212+
- deltaDataStart;
204213
long newValueCapacity = targetDataSize + deltaDataSize;
205214

206215
// make sure there is enough capacity
@@ -221,7 +230,7 @@ public ValueVector visit(BaseLargeVariableWidthVector deltaVector, Void value) {
221230

222231
// append data buffer
223232
MemoryUtil.copyMemory(
224-
deltaVector.getDataBuffer().memoryAddress(),
233+
deltaVector.getDataBuffer().memoryAddress() + deltaDataStart,
225234
targetVector.getDataBuffer().memoryAddress() + targetDataSize,
226235
deltaDataSize);
227236

@@ -232,7 +241,7 @@ public ValueVector visit(BaseLargeVariableWidthVector deltaVector, Void value) {
232241
+ (targetVector.getValueCount() + 1) * BaseLargeVariableWidthVector.OFFSET_WIDTH,
233242
deltaVector.getValueCount() * BaseLargeVariableWidthVector.OFFSET_WIDTH);
234243

235-
// increase each offset from the second buffer
244+
// rebase each appended offset to the target's data, accounting for the delta's start offset
236245
for (int i = 0; i < deltaVector.getValueCount(); i++) {
237246
long oldOffset =
238247
targetVector
@@ -245,7 +254,7 @@ public ValueVector visit(BaseLargeVariableWidthVector deltaVector, Void value) {
245254
.setLong(
246255
(long) (targetVector.getValueCount() + 1 + i)
247256
* BaseLargeVariableWidthVector.OFFSET_WIDTH,
248-
oldOffset + targetDataSize);
257+
oldOffset - deltaDataStart + targetDataSize);
249258
}
250259
((BaseLargeVariableWidthVector) targetVector).setLastSet(newValueCount - 1);
251260
targetVector.setValueCount(newValueCount);
@@ -331,16 +340,20 @@ public ValueVector visit(ListVector deltaVector, Void value) {
331340
targetVector
332341
.getOffsetBuffer()
333342
.getInt((long) targetVector.getValueCount() * ListVector.OFFSET_WIDTH);
334-
int deltaListSize =
343+
// see the corresponding comment in visit(BaseVariableWidthVector, Void): the delta's
344+
// offset buffer need not start at zero
345+
int deltaListStart = deltaVector.getOffsetBuffer().getInt(0);
346+
int deltaListEnd =
335347
deltaVector
336348
.getOffsetBuffer()
337349
.getInt((long) deltaVector.getValueCount() * ListVector.OFFSET_WIDTH);
350+
int deltaListSize = deltaListEnd - deltaListStart;
338351

339352
ListVector targetListVector = (ListVector) targetVector;
340353

341354
// make sure the underlying vector has value count set
342355
targetListVector.getDataVector().setValueCount(targetListSize);
343-
deltaVector.getDataVector().setValueCount(deltaListSize);
356+
deltaVector.getDataVector().setValueCount(deltaListEnd);
344357

345358
// make sure there is enough capacity
346359
while (targetVector.getValueCapacity() < newValueCount) {
@@ -372,13 +385,16 @@ public ValueVector visit(ListVector deltaVector, Void value) {
372385
.getOffsetBuffer()
373386
.setInt(
374387
(long) (targetVector.getValueCount() + 1 + i) * ListVector.OFFSET_WIDTH,
375-
oldOffset + targetListSize);
388+
oldOffset - deltaListStart + targetListSize);
376389
}
377390
targetListVector.setLastSet(newValueCount - 1);
378391

379392
// append underlying vectors
380-
VectorAppender innerAppender = new VectorAppender(targetListVector.getDataVector());
381-
deltaVector.getDataVector().accept(innerAppender, null);
393+
appendDataVector(
394+
targetListVector.getDataVector(),
395+
deltaVector.getDataVector(),
396+
deltaListStart,
397+
deltaListSize);
382398

383399
targetVector.setValueCount(newValueCount);
384400
return targetVector;
@@ -400,17 +416,21 @@ public ValueVector visit(LargeListVector deltaVector, Void value) {
400416
targetVector
401417
.getOffsetBuffer()
402418
.getLong((long) targetVector.getValueCount() * LargeListVector.OFFSET_WIDTH);
403-
long deltaListSize =
419+
// see the corresponding comment in visit(BaseVariableWidthVector, Void): the delta's
420+
// offset buffer need not start at zero
421+
long deltaListStart = deltaVector.getOffsetBuffer().getLong(0);
422+
long deltaListEnd =
404423
deltaVector
405424
.getOffsetBuffer()
406425
.getLong((long) deltaVector.getValueCount() * LargeListVector.OFFSET_WIDTH);
426+
long deltaListSize = deltaListEnd - deltaListStart;
407427

408-
ListVector targetListVector = (ListVector) targetVector;
428+
LargeListVector targetListVector = (LargeListVector) targetVector;
409429

410430
// make sure the underlying vector has value count set
411431
// todo recheck these casts when int64 vectors are supported
412432
targetListVector.getDataVector().setValueCount(checkedCastToInt(targetListSize));
413-
deltaVector.getDataVector().setValueCount(checkedCastToInt(deltaListSize));
433+
deltaVector.getDataVector().setValueCount(checkedCastToInt(deltaListEnd));
414434

415435
// make sure there is enough capacity
416436
while (targetVector.getValueCapacity() < newValueCount) {
@@ -427,10 +447,10 @@ public ValueVector visit(LargeListVector deltaVector, Void value) {
427447

428448
// append offset buffer
429449
MemoryUtil.copyMemory(
430-
deltaVector.getOffsetBuffer().memoryAddress() + ListVector.OFFSET_WIDTH,
450+
deltaVector.getOffsetBuffer().memoryAddress() + LargeListVector.OFFSET_WIDTH,
431451
targetVector.getOffsetBuffer().memoryAddress()
432452
+ (targetVector.getValueCount() + 1) * LargeListVector.OFFSET_WIDTH,
433-
(long) deltaVector.getValueCount() * ListVector.OFFSET_WIDTH);
453+
(long) deltaVector.getValueCount() * LargeListVector.OFFSET_WIDTH);
434454

435455
// increase each offset from the second buffer
436456
for (int i = 0; i < deltaVector.getValueCount(); i++) {
@@ -443,18 +463,42 @@ public ValueVector visit(LargeListVector deltaVector, Void value) {
443463
.getOffsetBuffer()
444464
.setLong(
445465
(long) (targetVector.getValueCount() + 1 + i) * LargeListVector.OFFSET_WIDTH,
446-
oldOffset + targetListSize);
466+
oldOffset - deltaListStart + targetListSize);
447467
}
448468
targetListVector.setLastSet(newValueCount - 1);
449469

450470
// append underlying vectors
451-
VectorAppender innerAppender = new VectorAppender(targetListVector.getDataVector());
452-
deltaVector.getDataVector().accept(innerAppender, null);
471+
appendDataVector(
472+
targetListVector.getDataVector(),
473+
deltaVector.getDataVector(),
474+
checkedCastToInt(deltaListStart),
475+
checkedCastToInt(deltaListSize));
453476

454477
targetVector.setValueCount(newValueCount);
455478
return targetVector;
456479
}
457480

481+
/**
482+
* Appends the range [start, start + length) of the delta vector's data vector to the target
483+
* vector's data vector. The range may not cover the whole delta data vector when the delta's
484+
* offset buffer does not start at zero.
485+
*/
486+
private static void appendDataVector(
487+
ValueVector targetDataVector, ValueVector deltaDataVector, int start, int length) {
488+
if (start == 0 && length == deltaDataVector.getValueCount()) {
489+
VectorAppender innerAppender = new VectorAppender(targetDataVector);
490+
deltaDataVector.accept(innerAppender, null);
491+
return;
492+
}
493+
TransferPair transferPair =
494+
deltaDataVector.getTransferPair(deltaDataVector.getField(), deltaDataVector.getAllocator());
495+
transferPair.splitAndTransfer(start, length);
496+
try (ValueVector slicedDeltaDataVector = transferPair.getTo()) {
497+
VectorAppender innerAppender = new VectorAppender(targetDataVector);
498+
slicedDeltaDataVector.accept(innerAppender, null);
499+
}
500+
}
501+
458502
@Override
459503
public ValueVector visit(FixedSizeListVector deltaVector, Void value) {
460504
Preconditions.checkArgument(

0 commit comments

Comments
 (0)