This is an automated email from the ASF dual-hosted git repository.

SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git


The following commit(s) were added to refs/heads/main by this push:
     new 3663237  fix(arrow): retain pools for stream buffers (#182)
3663237 is described below

commit 366323740cc13b5e113d9465080624b28e3880e8
Author: Mr Dk. <[email protected]>
AuthorDate: Thu Aug 6 18:26:25 2026 +0800

    fix(arrow): retain pools for stream buffers (#182)
    
    Keep Arrow memory pools alive while returned stream buffers remain
    referenced and note the allocation overhead for a future performance
    improvement.
    
    Add deterministic coverage for synchronous and asynchronous reads after
    adapter teardown.
    
    Co-authored-by: GPT-5.6 Terra <[email protected]>
---
 .../utils/arrow/arrow_input_stream_adapter.cpp     |  19 +++-
 .../utils/arrow/arrow_stream_adapter_test.cpp      | 120 +++++++++++++++++++++
 2 files changed, 137 insertions(+), 2 deletions(-)

diff --git a/src/paimon/common/utils/arrow/arrow_input_stream_adapter.cpp 
b/src/paimon/common/utils/arrow/arrow_input_stream_adapter.cpp
index 01429df..8c94027 100644
--- a/src/paimon/common/utils/arrow/arrow_input_stream_adapter.cpp
+++ b/src/paimon/common/utils/arrow/arrow_input_stream_adapter.cpp
@@ -43,6 +43,20 @@ arrow::Status ValidateArrowIoRange(int64_t value, const 
char* name) {
     return arrow::Status::OK();
 }
 
+struct BufferWithMemoryPool {
+    std::shared_ptr<arrow::MemoryPool> pool;
+    std::shared_ptr<arrow::Buffer> buffer;
+};
+
+std::shared_ptr<arrow::Buffer> 
KeepMemoryPoolAlive(std::shared_ptr<arrow::Buffer> buffer,
+                                                   const 
std::shared_ptr<arrow::MemoryPool>& pool) {
+    // TODO(lxy): Optimize the extra allocation introduced to retain the 
memory pool.
+    auto holder =
+        std::make_shared<BufferWithMemoryPool>(BufferWithMemoryPool{pool, 
std::move(buffer)});
+    auto* buffer_ptr = holder->buffer.get();
+    return std::shared_ptr<arrow::Buffer>(std::move(holder), buffer_ptr);
+}
+
 }  // namespace
 
 ArrowInputStreamAdapter::ArrowInputStreamAdapter(
@@ -82,7 +96,7 @@ arrow::Result<std::shared_ptr<arrow::Buffer>> 
ArrowInputStreamAdapter::Read(int6
     if (read_bytes < nbytes) {
         ARROW_RETURN_NOT_OK(buffer->Resize(read_bytes));
     }
-    return std::shared_ptr<arrow::Buffer>(std::move(buffer));
+    return 
KeepMemoryPoolAlive(std::shared_ptr<arrow::Buffer>(std::move(buffer)), pool_);
 }
 
 arrow::Result<int64_t> ArrowInputStreamAdapter::ReadAt(int64_t position, 
int64_t nbytes,
@@ -107,7 +121,7 @@ arrow::Result<std::shared_ptr<arrow::Buffer>> 
ArrowInputStreamAdapter::ReadAt(in
     if (read_bytes < nbytes) {
         ARROW_RETURN_NOT_OK(buffer->Resize(read_bytes));
     }
-    return std::shared_ptr<arrow::Buffer>(std::move(buffer));
+    return 
KeepMemoryPoolAlive(std::shared_ptr<arrow::Buffer>(std::move(buffer)), pool_);
 }
 
 arrow::Future<std::shared_ptr<arrow::Buffer>> 
ArrowInputStreamAdapter::ReadAsync(
@@ -131,6 +145,7 @@ arrow::Future<std::shared_ptr<arrow::Buffer>> 
ArrowInputStreamAdapter::ReadAsync
         return fut;
     }
     std::shared_ptr<arrow::Buffer> buffer = 
std::move(buffer_result).ValueUnsafe();
+    buffer = KeepMemoryPoolAlive(std::move(buffer), pool_);
     std::shared_ptr<std::atomic<uint64_t>> storage_read_bytes = 
storage_read_bytes_;
     input_stream_->ReadAsync(
         reinterpret_cast<char*>(buffer->mutable_data()), nbytes, position,
diff --git a/src/paimon/common/utils/arrow/arrow_stream_adapter_test.cpp 
b/src/paimon/common/utils/arrow/arrow_stream_adapter_test.cpp
index b568060..7280ad4 100644
--- a/src/paimon/common/utils/arrow/arrow_stream_adapter_test.cpp
+++ b/src/paimon/common/utils/arrow/arrow_stream_adapter_test.cpp
@@ -18,8 +18,11 @@
  */
 
 #include <cstdint>
+#include <cstring>
+#include <functional>
 #include <memory>
 #include <string>
+#include <utility>
 
 #include "arrow/api.h"
 #include "arrow/io/type_fwd.h"
@@ -35,6 +38,79 @@
 
 namespace paimon::test {
 
+namespace {
+
+constexpr char kTestPayload[] = "data";
+constexpr int64_t kTestSize = sizeof(kTestPayload) - 1;
+
+class DeferredInputStream : public InputStream {
+ public:
+    Status Seek(int64_t, SeekOrigin) override {
+        return Status::OK();
+    }
+
+    Result<int64_t> GetPos() const override {
+        return 0;
+    }
+
+    Result<int64_t> Read(char* buffer, int64_t size) override {
+        if (size != kTestSize) {
+            return Status::Invalid("unexpected read size");
+        }
+        std::memcpy(buffer, kTestPayload, kTestSize);
+        return kTestSize;
+    }
+
+    Result<int64_t> Read(char* buffer, int64_t size, int64_t) override {
+        return Read(buffer, size);
+    }
+
+    void ReadAsync(char* buffer, int64_t size, int64_t,
+                   std::function<void(Status)>&& callback) override {
+        buffer_ = buffer;
+        size_ = size;
+        callback_ = std::move(callback);
+    }
+
+    Status Complete() {
+        if (!callback_) {
+            return Status::Invalid("async request was not started");
+        }
+        if (size_ != kTestSize) {
+            return Status::Invalid("unexpected async read size");
+        }
+        std::memcpy(buffer_, kTestPayload, kTestSize);
+        auto callback = std::move(callback_);
+        callback(Status::OK());
+        return Status::OK();
+    }
+
+    Status Close() override {
+        return Status::OK();
+    }
+
+    Result<std::string> GetUri() const override {
+        return std::string("test://input");
+    }
+
+    Result<int64_t> Length() const override {
+        return kTestSize;
+    }
+
+ private:
+    char* buffer_ = nullptr;
+    int64_t size_ = 0;
+    std::function<void(Status)> callback_;
+};
+
+std::shared_ptr<ArrowInputStreamAdapter> CreateAdapter(
+    const std::shared_ptr<DeferredInputStream>& stream,
+    const std::shared_ptr<arrow::MemoryPool>& pool) {
+    return std::make_shared<ArrowInputStreamAdapter>(stream, kTestSize, pool);
+}
+
+}  // namespace
+
 TEST(ArrowStreamAdapterTest, TestInputAndOutputStream) {
     auto test_root_dir = UniqueTestDirectory::Create();
     ASSERT_TRUE(test_root_dir);
@@ -91,4 +167,48 @@ TEST(ArrowStreamAdapterTest, TestInputAndOutputStream) {
     ASSERT_TRUE(in_stream->closed());
 }
 
+TEST(ArrowStreamAdapterTest, TestReadKeepsMemoryPoolAliveUntilBufferReleased) {
+    auto stream = std::make_shared<DeferredInputStream>();
+    std::shared_ptr<arrow::MemoryPool> pool(GetArrowPool(GetDefaultPool()));
+    auto adapter = CreateAdapter(stream, pool);
+    ASSERT_NE(adapter, nullptr);
+
+    std::shared_ptr<arrow::Buffer> buffer = 
adapter->Read(kTestSize).ValueOrDie();
+    adapter.reset();
+    pool.reset();
+
+    ASSERT_EQ(buffer->ToString(), kTestPayload);
+    buffer.reset();
+}
+
+TEST(ArrowStreamAdapterTest, 
TestReadAtKeepsMemoryPoolAliveUntilBufferReleased) {
+    auto stream = std::make_shared<DeferredInputStream>();
+    std::shared_ptr<arrow::MemoryPool> pool(GetArrowPool(GetDefaultPool()));
+    auto adapter = CreateAdapter(stream, pool);
+    ASSERT_NE(adapter, nullptr);
+
+    std::shared_ptr<arrow::Buffer> buffer = adapter->ReadAt(0, 
kTestSize).ValueOrDie();
+    adapter.reset();
+    pool.reset();
+
+    ASSERT_EQ(buffer->ToString(), kTestPayload);
+    buffer.reset();
+}
+
+TEST(ArrowStreamAdapterTest, 
TestAsyncReadKeepsMemoryPoolAliveUntilBufferReleased) {
+    auto stream = std::make_shared<DeferredInputStream>();
+    std::shared_ptr<arrow::MemoryPool> pool(GetArrowPool(GetDefaultPool()));
+    auto adapter = CreateAdapter(stream, pool);
+    ASSERT_NE(adapter, nullptr);
+
+    auto future = adapter->ReadAsync(arrow::io::default_io_context(), 0, 
kTestSize);
+    adapter.reset();
+    pool.reset();
+
+    ASSERT_OK(stream->Complete());
+    std::shared_ptr<arrow::Buffer> buffer = future.MoveResult().ValueOrDie();
+    ASSERT_EQ(buffer->ToString(), kTestPayload);
+    buffer.reset();
+}
+
 }  // namespace paimon::test

Reply via email to