feat(ipc): support dictionary updates and deltas - #928
rustyconover wants to merge 16 commits into
Conversation
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #928 +/- ##
==========================================
- Coverage 78.16% 77.97% -0.20%
==========================================
Files 106 106
Lines 16842 17154 +312
Branches 1984 2067 +83
==========================================
+ Hits 13165 13376 +211
- Misses 2443 2502 +59
- Partials 1234 1276 +42 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
0e55dbe to
de091e0
Compare
paleolimbot
left a comment
There was a problem hiding this comment.
Thank you for this!
I will take a closer look in the next few days, but can I ask that the delta dictionary and/or appender piece be separated from the IPC write piece? All three are independently useful and I think it will be a smoother process if they're PRed separately as well.
|
Yep I'll break them out! Rusty |
de091e0 to
8db5007
Compare
|
Split completed as requested:
The original full branch is preserved as |
8db5007 to
0951095
Compare
0951095 to
067826a
Compare
## Summary - add `ArrowArrayAppendArrayView()` as a core array-building API - append primitive, nested, union, list-view, and run-end encoded values from an `ArrowArrayView` - preserve source slicing semantics and report unsupported sliced run-end encoded inputs - add direct native coverage for primitive, nested/sliced, union, and run-end encoded arrays ## Context This extracts the generic array-view appender from #928 as requested in review. It is independently useful and will be the prerequisite for a separate dictionary-delta decoding PR; #928 no longer contains the appender or decoder changes. ## Validation - native CMake build passed - native CTest suite: 334/334 passed (4 optional codec tests skipped)
067826a to
6b8b6c7
Compare
paleolimbot
left a comment
There was a problem hiding this comment.
Thank you for this!
I think this PR contains 3 orthogonal changes:
- Decoding delta dictionaries + some small fixes to the decoder
- Avoiding dictionary replacement when emitting dictionaries that are identical + assorted fixes to the encoder for nested dictionaries
- Encoding/writing delta dictionaries
I did a first pass on everything and I think we can get all these merged soon!
| static void ArrowIpcArrayPrepareForAppend(struct ArrowArray* array, | ||
| const struct ArrowArrayView* array_view) { | ||
| // Finishing a view array materializes its variadic-buffer sizes. Appending may | ||
| // extend the last variadic buffer or add another one, so force the sizes buffer | ||
| // to be regenerated by the next ArrowArrayFinishBuildingDefault(). | ||
| if (array_view->storage_type == NANOARROW_TYPE_BINARY_VIEW || | ||
| array_view->storage_type == NANOARROW_TYPE_STRING_VIEW) { | ||
| ArrowBufferReset(ArrowArrayBuffer(array, array->n_buffers - 1)); | ||
| } | ||
|
|
||
| for (int64_t i = 0; i < array->n_children; i++) { | ||
| ArrowIpcArrayPrepareForAppend(array->children[i], array_view->children[i]); | ||
| } | ||
| } |
There was a problem hiding this comment.
I don't think we support views yet for IPC decoding at all yet? #824
| ArrowArrayMove(&dictionary->current_value, &combined); | ||
|
|
||
| ArrowIpcArrayPrepareForAppend(&combined, array_view); | ||
| ArrowErrorCode result = ArrowArrayReserve(&combined, value->length); | ||
| if (result == NANOARROW_OK) { | ||
| result = ArrowArrayAppendStorageFromArrayView(&combined, array_view, error); | ||
| } | ||
| if (result == NANOARROW_OK) { | ||
| result = ArrowIpcArraySetDictionaries(&combined, value); | ||
| } |
There was a problem hiding this comment.
For the PR that does contain this, a good pattern is to define a function that can fail and doesn't own any input pointers:
ArrowErrorCode ArrowIpcDoAppend(struct ArrowArray* combined, struct ArrowArrayView array_view, struct ArrowError* error) {
NANOARROW_RETURN_NOT_OK(...);
return NANOARROW_OK;
}then here you can just check once and release the temporary
ArrowErrorCode result = ArrowIpcDoAppend(...);
if (result != NANOARROW_OK) {
// release stuff
return result;
}| enum class DeltaDictionaryValueCase { | ||
| kBoolean, | ||
| kInt64, | ||
| kInt64WithNull, | ||
| kUInt64, | ||
| kDouble, | ||
| kString, | ||
| kBinary, | ||
| kDecimal128, | ||
| kList, | ||
| kStruct, | ||
| kFixedSizeList | ||
| }; |
There was a problem hiding this comment.
Can you reuse enum ArrowType for this one?
| TEST(NanoarrowIpcTest, RejectsArrowCppDenseUnionDictionaryDelta) { | ||
| arrow::Int8Builder type_ids_builder; | ||
| arrow::Int32Builder offsets_builder; |
There was a problem hiding this comment.
I think we can skip this test unless there's something special about rejecting a type in the decoder implementation
| // In the usual streaming loop, the previously returned batch has been released | ||
| // before the next one is requested. Recover the mutable backing array and append | ||
| // directly so a sequence of small deltas grows geometrically instead of copying | ||
| // the complete dictionary for every message. If an older batch is still alive, | ||
| // keep the copy-on-write path below to preserve its dictionary snapshot. | ||
| if (ArrowArrayInternalTryUnshare(&dictionary->current_value)) { | ||
| struct ArrowArray combined; | ||
| ArrowArrayMove(&dictionary->current_value, &combined); |
There was a problem hiding this comment.
The usual streaming loop should be using ArrowIpcDecoderDecodeArrayViewWithDictionaries(), which in a perfect world never was shared (however, it probably had a shared backing buffer on its first decode because most decoding happens from a shared buffer).
Rather than trying to unshare, can we check whether this array can be appended to (and if possibly delay the sharing of it until it is requested as an array that is not just a view?). Totally ok if not possible, but if so, we should punt on the trying to unshare piece and implement that in a follow up.
| static void ArrowIpcWriterCanonicalizeBitmapPadding( | ||
| struct ArrowArray* array, const struct ArrowArrayView* array_view) { | ||
| int64_t remainder = array->length % 8; |
There was a problem hiding this comment.
It would help to put the comment about canonicalizing the padding helping to minimize reemitting a dictionary here.
Other parts of nanoarrow typically avoid this by zeroing out the last byte of a bitmap when reserving but it is sometimes hard to guarantee that everywhere.
| ArrowErrorCode result = ArrowArrayInitFromArrayView(out, src, error); | ||
| if (result == NANOARROW_OK) { | ||
| result = ArrowArrayStartAppending(out); | ||
| } |
There was a problem hiding this comment.
The same nested function strategy for avoiding these repeated result checks applies here, too
| static ArrowErrorCode ArrowIpcWriterCompareMaterializedArrays( | ||
| const struct ArrowArray* lhs, const struct ArrowArray* rhs, | ||
| const struct ArrowArrayView* shape, int* out, struct ArrowError* error) { |
There was a problem hiding this comment.
Rather than compare by value (possibly slower than just writing another dictionary to the stream), probably comparing the buffer pointers makes more sense. If it's a shared array it should be pointing to the same data.
| if (result != NANOARROW_OK) { | ||
| goto cleanup; | ||
| } | ||
| is_delta = 1; |
There was a problem hiding this comment.
You can use the nested function approach for shared cleanup (we don't use goto in nanoarrow, or at least we haven't yet)
| if (allow_delta && !force_emit && cached != NULL && | ||
| current_values.length > cached->values.length) { | ||
| result = ArrowIpcWriterMaterializeArrayView(values_view, 0, cached->values.length, | ||
| &prefix_values, error); | ||
| if (result != NANOARROW_OK) { | ||
| goto cleanup; | ||
| } |
There was a problem hiding this comment.
I don't think that the writer should emit delta dictionaries in this way...this is a kind of calculation that application code should do, which may have access to things like C++ that can do it more efficiently.
In any case, emitting delta dictionaries is a third scope here for a separate PR
Summary
ArrowArrayAppendStorageFromArrayView()Relationship to #930
#930 has merged, and this branch has been rebased onto the resulting
main. This PR now uses the publicArrowArrayAppendStorageFromArrayView()API from #930; the array-view appender implementation and #930 commits are no longer part of this diff.The remaining changes in this PR are the dictionary IPC update and delta functionality.
IPC behavior
Streams cache dictionary values, omit unchanged dictionaries, emit append-only suffixes with
DictionaryBatch.isDelta, and advertiseDICTIONARY_REPLACEMENTwhen a full replacement may occur. Nested dictionaries are emitted dependency-first and dependents are invalidated when needed.Files permit append-only delta dictionaries but reject full dictionary replacement after a dictionary ID has been written.
Delta decoding supports the storage types handled by
ArrowArrayAppendStorageFromArrayView(). Dense union deltas currently returnENOTSUP, matching #930 support.Validation