Native Arrow path support in stream transport with zero-copy transfer - #21253
Conversation
PR Reviewer Guide 🔍(Review updated until commit a5e7857)Here are some key observations to aid the review process:
|
28617d9 to
7d506e8
Compare
PR Code Suggestions ✨Latest suggestions up to a5e7857 Explore these optional code suggestions:
Previous suggestionsSuggestions up to commit 2dce758
Suggestions up to commit af5cc9c
Suggestions up to commit 7f0c5f0
Suggestions up to commit 416abe6
Suggestions up to commit d65e0f6
|
|
Persistent review updated to latest commit 7d506e8 |
7d506e8 to
ac83ec6
Compare
|
Persistent review updated to latest commit ac83ec6 |
|
❌ Gradle check result for ac83ec6: FAILURE Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change? |
ac83ec6 to
e3d1a5e
Compare
|
Persistent review updated to latest commit e3d1a5e |
e3d1a5e to
ac213cc
Compare
|
Persistent review updated to latest commit ac213cc |
ac213cc to
9280e70
Compare
|
Persistent review updated to latest commit 9280e70 |
9280e70 to
3118626
Compare
|
Persistent review updated to latest commit 3118626 |
3118626 to
73d7245
Compare
|
Persistent review updated to latest commit 73d7245 |
73d7245 to
d76afed
Compare
|
Persistent review updated to latest commit d76afed |
|
❌ Gradle check result for d76afed: FAILURE Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change? |
d76afed to
52d8ab3
Compare
|
Persistent review updated to latest commit 52d8ab3 |
52d8ab3 to
c9e78c5
Compare
|
Persistent review updated to latest commit c9e78c5 |
|
❌ Gradle check result for c9e78c5: null Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change? |
c9e78c5 to
c037252
Compare
|
Persistent review updated to latest commit c037252 |
|
❌ Gradle check result for c037252: FAILURE Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change? |
…ct#21253) Signed-off-by: Rishabh Maurya <rishabhmaurya05@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
Today's path drains DataFusion results into Object[] rows, sends one buffered response, and the coordinator converts back to Arrow. Replace with native Arrow batches over the stream transport from opensearch-project#21253. Signed-off-by: bowenlan-amzn <bowenlan23@gmail.com>
…ct#21253) Signed-off-by: Rishabh Maurya <rishabhmaurya05@gmail.com>
Context
Adds a native Arrow transport path to the Flight transport plugin, enabling zero-copy transfer of Arrow data without byte serialization. This is an alternative to #21240 that keeps all changes within the
arrow-flight-rpcplugin.Problem
The existing byte-serialized path (
VectorStreamOutput.ByteSerialized) serializes Arrow data into aVarBinaryVector, which is then sent via FlightputNext(). The byte-serialized path is essential for streaming aggregation where results are produced as OpenSearch objects and need serialization into Arrow format. However, for use cases like DataFusion integration where data already originates as native Arrow vectors (potentially from C data import), a direct transfer path avoids the serialization round-trip.Solution
Introduce
ArrowBatchResponse— an abstract base class that API developers extend. When the framework detects this response type, it performs a zero-copytransferTo()of the producer's vectors into the channel's shared root, bypassing serialization entirely.Key design decisions
Producer's allocator for the shared root (same-allocator transfer). The shared root is created from the first batch's producer allocator. This ensures same-allocator transfer, which avoids an Arrow Java bug where
BufferLedger.transferOwnership()of foreign-backed buffers (from C data import viawrapForeignAllocation) doesn't properly free theArrowArrayC struct (128 bytes per batch). Same-allocator transfer sidesteps this entirely.Long-lived allocator requirement. The allocator used for producer roots must outlive the gRPC stream. gRPC's zero-copy write path (
ArrowBufRetainingCompositeByteBuf) retains ArrowBuf references beyondputNext()and even beyondcompleted()— they are released asynchronously by gRPC's Netty event loop. Closing the allocator while gRPC still holds these retained references causes memory accounting errors.Producer root closed after transfer. Each batch's producer root is closed by the framework after
transferTo()moves its buffers into the shared root. The producer's buffers are empty after transfer, so close is safe and immediate.Challenges investigated
gRPC zero-copy buffer lifecycle.
putNext()withsetUseZeroCopy(true)createsArrowBufRetainingCompositeByteBufwhich retains ArrowBufs independently of the shared root. WhentransferTo()replaces the shared root's buffers on the next batch, the old ArrowBufs are released from the root's side but kept alive by gRPC's retain (refcount > 0). The byte-serialized path avoids this because it reuses the same ArrowBuf across batches (overwriting contents viasetSafe()), so gRPC's retained ByteBufs always point to valid, live memory. For the native arrow path, the allocator must be long-lived so gRPC can release the retained ArrowBufs back to it at any time.Upstream issues discovered
During this work we identified two issues in Arrow Java that are worth reporting upstream:
Arrow Java:
BufferLedger.transferOwnership()leaksArrowArrayC struct for foreign-backed buffers. WhenArrayImporter.importArray()imports data via the C Data Interface, it allocates a 128-byteArrowArrayC struct buffer from the data allocator, wrapped inReferenceCountedArrowArray. During cross-allocatortransferOwnership, the accounting for the data buffers moves correctly, but theReferenceCountedArrowArrayrefcount never reaches 0 because the foreign allocation cleanup path isn't triggered. This leaks 128 bytes per imported array per cross-allocator transfer. Same-allocator transfer avoids this becausetransferOwnershipreturns the same buffer (no new allocation). Present in Arrow Java 18.1.0 and latest main.Arrow Flight Java:
ServerStreamListener.completed()is fire-and-forget with no flush guarantee.completed()enqueues HTTP/2 trailers to gRPC'sWriteQueuebut returns immediately without waiting for pending data frames to be flushed. Meanwhile,ArrowBufRetainingCompositeByteBufholds retained references to ArrowBufs that are only released when gRPC's Netty event loop processes the write. There is no callback mechanism (ServerStreamListenerdoesn't exposesetOnCloseHandler, which exists onServerCallStreamObserverbut is not accessible through the Flight API) to know when gRPC has fully released all buffer references. This means allocators backing the stream's data cannot be safely closed immediately aftercompleted(). AsetOnCloseHandleror similar API onServerStreamListenerwould allow producers to defer cleanup until gRPC is truly done with the buffers.Design doc
See
plugins/arrow-flight-rpc/docs/native-arrow-transport-design.md