From 5a777841ad70d5b4e8d2055decffd27ba31c5194 Mon Sep 17 00:00:00 2001 From: Ziy1-Tan Date: Sun, 29 Jun 2025 16:38:18 +0800 Subject: [PATCH 1/2] [C++][Parquet] Add arrow::Result version of parquet::arrow::FileReader::GetRecordBatchReader() Signed-off-by: Ziy1-Tan --- cpp/examples/arrow/parquet_read_write.cc | 2 +- cpp/src/parquet/arrow/reader.cc | 41 +++++++++++------------- cpp/src/parquet/arrow/reader.h | 40 +++++++++++++++++++++++ 3 files changed, 60 insertions(+), 23 deletions(-) diff --git a/cpp/examples/arrow/parquet_read_write.cc b/cpp/examples/arrow/parquet_read_write.cc index 246501896637..9284ae895fc1 100644 --- a/cpp/examples/arrow/parquet_read_write.cc +++ b/cpp/examples/arrow/parquet_read_write.cc @@ -67,7 +67,7 @@ arrow::Status ReadInBatches(std::string path_to_file) { ARROW_ASSIGN_OR_RAISE(arrow_reader, reader_builder.Build()); std::shared_ptr<::arrow::RecordBatchReader> rb_reader; - ARROW_ASSIGN_OR_RAISE(rb_reader, arrow_reader->GetRecordBatchReader()); + ARROW_ASSIGN_OR_RAISE(rb_reader, arrow_reader->GetRecordBatchReaderSharedPtr()); for (arrow::Result> maybe_batch : *rb_reader) { // Operate on each batch... diff --git a/cpp/src/parquet/arrow/reader.cc b/cpp/src/parquet/arrow/reader.cc index d9a2b989148d..27c5d938b9a1 100644 --- a/cpp/src/parquet/arrow/reader.cc +++ b/cpp/src/parquet/arrow/reader.cc @@ -342,6 +342,25 @@ class FileReaderImpl : public FileReader { Iota(reader_->metadata()->num_columns())); } + Result> GetRecordBatchReaderSharedPtr() override { + ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader()); + return std::shared_ptr(std::move(tmp)); + } + + Result> GetRecordBatchReaderSharedPtr( + const std::vector& row_group_indices) override { + ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader(row_group_indices)); + return std::shared_ptr(std::move(tmp)); + } + + Result> GetRecordBatchReaderSharedPtr( + const std::vector& row_group_indices, + const std::vector& column_indices) override { + ARROW_ASSIGN_OR_RAISE(auto tmp, + GetRecordBatchReader(row_group_indices, column_indices)); + return std::shared_ptr(std::move(tmp)); + } + ::arrow::Result<::arrow::AsyncGenerator>> GetRecordBatchGenerator(std::shared_ptr reader, const std::vector row_group_indices, @@ -1310,28 +1329,6 @@ std::shared_ptr FileReaderImpl::RowGroup(int row_group_index) { // ---------------------------------------------------------------------- // Public factory functions -Status FileReader::GetRecordBatchReader(std::shared_ptr* out) { - ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader()); - out->reset(tmp.release()); - return Status::OK(); -} - -Status FileReader::GetRecordBatchReader(const std::vector& row_group_indices, - std::shared_ptr* out) { - ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader(row_group_indices)); - out->reset(tmp.release()); - return Status::OK(); -} - -Status FileReader::GetRecordBatchReader(const std::vector& row_group_indices, - const std::vector& column_indices, - std::shared_ptr* out) { - ARROW_ASSIGN_OR_RAISE(auto tmp, - GetRecordBatchReader(row_group_indices, column_indices)); - out->reset(tmp.release()); - return Status::OK(); -} - Status FileReader::Make(::arrow::MemoryPool* pool, std::unique_ptr reader, const ArrowReaderProperties& properties, diff --git a/cpp/src/parquet/arrow/reader.h b/cpp/src/parquet/arrow/reader.h index 4a01d7c4e4bc..9b59dd9c2fac 100644 --- a/cpp/src/parquet/arrow/reader.h +++ b/cpp/src/parquet/arrow/reader.h @@ -191,13 +191,53 @@ class PARQUET_EXPORT FileReader { /// /// \returns error Status if either row_group_indices or column_indices /// contains an invalid index + /// \deprecated Deprecated in future release. Use arrow::Result version instead. + ARROW_DEPRECATED("Deprecated in future release. Use arrow::Result version instead.") ::arrow::Status GetRecordBatchReader(const std::vector& row_group_indices, const std::vector& column_indices, std::shared_ptr<::arrow::RecordBatchReader>* out); + + /// \brief Return a RecordBatchReader of row groups selected from + /// row_group_indices, whose columns are selected by column_indices. + /// + /// Note that the ordering in row_group_indices and column_indices + /// matter. FileReaders must outlive their RecordBatchReaders. + /// + /// \param row_group_indices which row groups to read (order determines read order). + /// \param column_indices which columns to read (order determines output schema). + /// + /// \returns error Result if either row_group_indices or column_indices + /// contains an invalid index + virtual ::arrow::Result> + GetRecordBatchReaderSharedPtr(const std::vector& row_group_indices, + const std::vector& column_indices) = 0; + + /// \deprecated Deprecated in future release. Use arrow::Result version instead. + ARROW_DEPRECATED("Deprecated in future release. Use arrow::Result version instead.") ::arrow::Status GetRecordBatchReader(const std::vector& row_group_indices, std::shared_ptr<::arrow::RecordBatchReader>* out); + + /// \brief Return a RecordBatchReader of row groups selected from row_group_indices. + /// + /// Note that the ordering in row_group_indices matters. FileReaders must outlive + /// their RecordBatchReaders. + /// + /// \param row_group_indices which row groups to read (order determines read order). + /// + /// \returns error Result if row_group_indices contains an invalid index + virtual ::arrow::Result> + GetRecordBatchReaderSharedPtr(const std::vector& row_group_indices) = 0; + + /// \deprecated Deprecated in future release. Use arrow::Result version instead. + ARROW_DEPRECATED("Deprecated in future release. Use arrow::Result version instead.") ::arrow::Status GetRecordBatchReader(std::shared_ptr<::arrow::RecordBatchReader>* out); + /// \brief Return a RecordBatchReader of all row groups and columns. + /// + /// \returns error Result if row_group_indices contains an invalid index + virtual ::arrow::Result> + GetRecordBatchReaderSharedPtr() = 0; + /// \brief Return a generator of record batches. /// /// The FileReader must outlive the generator, so this requires that you pass in a From f285b941badf367b37b70ce5504066cb92e64de1 Mon Sep 17 00:00:00 2001 From: Ziy1-Tan Date: Tue, 1 Jul 2025 23:26:40 +0800 Subject: [PATCH 2/2] Deprecate shared_ptr version Signed-off-by: Ziy1-Tan --- cpp/examples/arrow/parquet_read_write.cc | 2 +- cpp/src/parquet/arrow/reader.cc | 41 ++++++++++++---------- cpp/src/parquet/arrow/reader.h | 44 ++++-------------------- 3 files changed, 29 insertions(+), 58 deletions(-) diff --git a/cpp/examples/arrow/parquet_read_write.cc b/cpp/examples/arrow/parquet_read_write.cc index 9284ae895fc1..246501896637 100644 --- a/cpp/examples/arrow/parquet_read_write.cc +++ b/cpp/examples/arrow/parquet_read_write.cc @@ -67,7 +67,7 @@ arrow::Status ReadInBatches(std::string path_to_file) { ARROW_ASSIGN_OR_RAISE(arrow_reader, reader_builder.Build()); std::shared_ptr<::arrow::RecordBatchReader> rb_reader; - ARROW_ASSIGN_OR_RAISE(rb_reader, arrow_reader->GetRecordBatchReaderSharedPtr()); + ARROW_ASSIGN_OR_RAISE(rb_reader, arrow_reader->GetRecordBatchReader()); for (arrow::Result> maybe_batch : *rb_reader) { // Operate on each batch... diff --git a/cpp/src/parquet/arrow/reader.cc b/cpp/src/parquet/arrow/reader.cc index 27c5d938b9a1..d9a2b989148d 100644 --- a/cpp/src/parquet/arrow/reader.cc +++ b/cpp/src/parquet/arrow/reader.cc @@ -342,25 +342,6 @@ class FileReaderImpl : public FileReader { Iota(reader_->metadata()->num_columns())); } - Result> GetRecordBatchReaderSharedPtr() override { - ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader()); - return std::shared_ptr(std::move(tmp)); - } - - Result> GetRecordBatchReaderSharedPtr( - const std::vector& row_group_indices) override { - ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader(row_group_indices)); - return std::shared_ptr(std::move(tmp)); - } - - Result> GetRecordBatchReaderSharedPtr( - const std::vector& row_group_indices, - const std::vector& column_indices) override { - ARROW_ASSIGN_OR_RAISE(auto tmp, - GetRecordBatchReader(row_group_indices, column_indices)); - return std::shared_ptr(std::move(tmp)); - } - ::arrow::Result<::arrow::AsyncGenerator>> GetRecordBatchGenerator(std::shared_ptr reader, const std::vector row_group_indices, @@ -1329,6 +1310,28 @@ std::shared_ptr FileReaderImpl::RowGroup(int row_group_index) { // ---------------------------------------------------------------------- // Public factory functions +Status FileReader::GetRecordBatchReader(std::shared_ptr* out) { + ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader()); + out->reset(tmp.release()); + return Status::OK(); +} + +Status FileReader::GetRecordBatchReader(const std::vector& row_group_indices, + std::shared_ptr* out) { + ARROW_ASSIGN_OR_RAISE(auto tmp, GetRecordBatchReader(row_group_indices)); + out->reset(tmp.release()); + return Status::OK(); +} + +Status FileReader::GetRecordBatchReader(const std::vector& row_group_indices, + const std::vector& column_indices, + std::shared_ptr* out) { + ARROW_ASSIGN_OR_RAISE(auto tmp, + GetRecordBatchReader(row_group_indices, column_indices)); + out->reset(tmp.release()); + return Status::OK(); +} + Status FileReader::Make(::arrow::MemoryPool* pool, std::unique_ptr reader, const ArrowReaderProperties& properties, diff --git a/cpp/src/parquet/arrow/reader.h b/cpp/src/parquet/arrow/reader.h index 9b59dd9c2fac..9753fe47ba06 100644 --- a/cpp/src/parquet/arrow/reader.h +++ b/cpp/src/parquet/arrow/reader.h @@ -191,53 +191,21 @@ class PARQUET_EXPORT FileReader { /// /// \returns error Status if either row_group_indices or column_indices /// contains an invalid index - /// \deprecated Deprecated in future release. Use arrow::Result version instead. - ARROW_DEPRECATED("Deprecated in future release. Use arrow::Result version instead.") + /// \deprecated Deprecated in 21.0.0. Use arrow::Result version instead. + ARROW_DEPRECATED("Deprecated in 21.0.0. Use arrow::Result version instead.") ::arrow::Status GetRecordBatchReader(const std::vector& row_group_indices, const std::vector& column_indices, std::shared_ptr<::arrow::RecordBatchReader>* out); - /// \brief Return a RecordBatchReader of row groups selected from - /// row_group_indices, whose columns are selected by column_indices. - /// - /// Note that the ordering in row_group_indices and column_indices - /// matter. FileReaders must outlive their RecordBatchReaders. - /// - /// \param row_group_indices which row groups to read (order determines read order). - /// \param column_indices which columns to read (order determines output schema). - /// - /// \returns error Result if either row_group_indices or column_indices - /// contains an invalid index - virtual ::arrow::Result> - GetRecordBatchReaderSharedPtr(const std::vector& row_group_indices, - const std::vector& column_indices) = 0; - - /// \deprecated Deprecated in future release. Use arrow::Result version instead. - ARROW_DEPRECATED("Deprecated in future release. Use arrow::Result version instead.") + /// \deprecated Deprecated in 21.0.0. Use arrow::Result version instead. + ARROW_DEPRECATED("Deprecated in 21.0.0. Use arrow::Result version instead.") ::arrow::Status GetRecordBatchReader(const std::vector& row_group_indices, std::shared_ptr<::arrow::RecordBatchReader>* out); - /// \brief Return a RecordBatchReader of row groups selected from row_group_indices. - /// - /// Note that the ordering in row_group_indices matters. FileReaders must outlive - /// their RecordBatchReaders. - /// - /// \param row_group_indices which row groups to read (order determines read order). - /// - /// \returns error Result if row_group_indices contains an invalid index - virtual ::arrow::Result> - GetRecordBatchReaderSharedPtr(const std::vector& row_group_indices) = 0; - - /// \deprecated Deprecated in future release. Use arrow::Result version instead. - ARROW_DEPRECATED("Deprecated in future release. Use arrow::Result version instead.") + /// \deprecated Deprecated in 21.0.0. Use arrow::Result version instead. + ARROW_DEPRECATED("Deprecated in 21.0.0. Use arrow::Result version instead.") ::arrow::Status GetRecordBatchReader(std::shared_ptr<::arrow::RecordBatchReader>* out); - /// \brief Return a RecordBatchReader of all row groups and columns. - /// - /// \returns error Result if row_group_indices contains an invalid index - virtual ::arrow::Result> - GetRecordBatchReaderSharedPtr() = 0; - /// \brief Return a generator of record batches. /// /// The FileReader must outlive the generator, so this requires that you pass in a