Copilot commented on code in PR #50629:
URL: https://github.com/apache/arrow/pull/50629#discussion_r4134546351


##########
cpp/src/parquet/column_reader.cc:
##########
@@ -724,357 +690,800 @@ class SkippableTypedDecoder {
   }
 };
 
-// ----------------------------------------------------------------------
-// Impl base class for TypedColumnReader and RecordReader
+/*********************
+ *  ValueSinkCursor  *
+ *********************/
 
-template <typename DType>
-class ColumnReaderImplBase {
+inline int64_t compute_capacity_pow2(int64_t capacity, int64_t size, int64_t 
extra_size) {
+  if (extra_size < 0) {
+    throw ParquetException("Negative size (corrupt file?)");
+  }
+  int64_t target_size = -1;
+  if (AddWithOverflow(size, extra_size, &target_size)) {
+    throw ParquetException("Allocation size too large (corrupt file?)");
+  }
+  if (target_size >= (1LL << 62)) {
+    throw ParquetException("Allocation size too large (corrupt file?)");
+  }
+  if (capacity >= target_size) {
+    return capacity;
+  }
+  return bit_util::NextPower2(target_size);
+}
+
+/// A simple cursor with a number of values and an available capacity.
+///
+/// This is a reused foundation for creating data sinks.
+/// Since decoders write to an already available buffer, we need to manually 
track a
+/// capacity where the user is allowed to write (contrary to say `std::vector` 
where
+/// it is UB to write in range `[size(), capacity()[`).
+class ValueSinkCursor {
  public:
-  using T = typename DType::c_type;
+  int64_t capacity() const { return capacity_; }
 
-  ColumnReaderImplBase(const ColumnDescriptor* descr, ::arrow::MemoryPool* 
pool)
-      : descr_(descr),
-        definition_level_decoder_(descr->max_definition_level()),
-        repetition_level_decoder_(descr->max_repetition_level()),
-        pool_(pool),
-        current_decoder_(pool) {}
+  int64_t values_count() const { return values_count_; }
 
-  virtual ~ColumnReaderImplBase() = default;
+  void set_values_count(int64_t vals) { values_count_ = vals; }
 
- protected:
-  // Read up to batch_size values from the current data page into the
-  // pre-allocated memory T*
-  //
-  // @returns: the number of values read into the out buffer
-  int64_t ReadValues(int64_t batch_size, T* out) {
-    int64_t num_decoded = current_decoder_->Decode(out, 
static_cast<int>(batch_size));
-    return num_decoded;
+  void increase_values_count(int64_t extra) { values_count_ += extra; }
+
+  int64_t fit_capacity_for_extra(int64_t extra_values) {
+    auto new_capacity = compute_capacity_pow2(capacity_, values_count_, 
extra_values);
+    ARROW_DCHECK_GE(new_capacity, capacity());
+    return std::exchange(capacity_, new_capacity);
   }
 
-  // Read up to batch_size values from the current data page into the
-  // pre-allocated memory T*, leaving spaces for null entries according
-  // to the def_levels.
-  //
-  // @returns: the number of values read into the out buffer
-  int64_t ReadValuesSpaced(int64_t batch_size, T* out, int64_t null_count,
-                           uint8_t* valid_bits, int64_t valid_bits_offset) {
-    return current_decoder_->DecodeSpaced(out, static_cast<int>(batch_size),
-                                          static_cast<int>(null_count), 
valid_bits,
-                                          valid_bits_offset);
+  void reset() {
+    values_count_ = 0;
+    capacity_ = 0;
   }
 
-  // Read multiple definition levels into preallocated memory
-  //
-  // Returns the number of decoded definition levels
-  int64_t ReadDefinitionLevels(int64_t batch_size, int16_t* levels) {
-    if (max_def_level() == 0) {
-      return 0;
+ private:
+  int64_t values_count_ = 0;
+  int64_t capacity_ = 0;
+};
+
+/********************
+ *  DataSinkBuffer  *
+ ********************/
+
+template <typename T>
+struct BytesCounterForValues {
+  /// Values are always fully written before being read, no need to zero new 
capacity.
+  static constexpr bool kZeroNewCapacity = false;
+
+  static int64_t bytes_for_values(int64_t nitems) {
+    constexpr auto kValueByteSize = static_cast<int64_t>(sizeof(T));
+    int64_t bytes = -1;
+    if (MultiplyWithOverflow(nitems, kValueByteSize, &bytes)) {
+      throw ParquetException("Total size of items too large");
     }
-    return definition_level_decoder_.Decode(static_cast<int>(batch_size), 
levels);
+    return bytes;
   }
+};
 
-  bool HasNextInternal() {
-    // Either there is no data page available yet, or the data page has been
-    // exhausted
-    if (num_buffered_values_ == 0 || num_decoded_values_ == 
num_buffered_values_) {
-      if (!ReadNewPage() || num_buffered_values_ == 0) {
-        return false;
-      }
+/// A base data sink that writes into an Arrow buffer.
+template <typename T, typename BytesCounter = BytesCounterForValues<T>>
+class DataSinkBuffer : private ValueSinkCursor {
+ public:
+  using value_type = T;
+
+  using ValueSinkCursor::capacity;
+  using ValueSinkCursor::values_count;
+
+  explicit DataSinkBuffer(std::shared_ptr<::arrow::ResizableBuffer> buffer)
+      : values_(std::move(buffer)) {}
+
+  bool is_void() const { return values_ == nullptr; }
+
+  value_type* data() const {
+    return is_void() ? nullptr : values_->mutable_data_as<value_type>();
+  }
+
+  /// Erase a range of element from the buffer.
+  void Erase(int64_t start, int64_t end) {
+    if (is_void() || start >= values_count() || start >= end) {
+      return;
     }
-    return true;
+    const auto last = std::min(end, values_count());
+    const auto count = last - start;
+    std::copy(data() + last, write_start(), data() + start);

Review Comment:
   `std::copy` is not permitted when the source and destination ranges overlap, 
but this erase path shifts the remaining levels left and normally overlaps them 
(for example, erasing the first level copies `[data + 1, end)` into `[data, 
...)`). That is undefined behavior and can corrupt buffered levels during 
repeated `SkipRecords`; use an overlap-safe move/copy operation instead.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to