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