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


##########
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);
+    set_values_count(values_count() - count);
   }
 
-  // Read multiple repetition levels into preallocated memory
-  // Returns the number of decoded repetition levels
-  int64_t ReadRepetitionLevels(int64_t batch_size, int16_t* levels) {
-    if (max_rep_level() == 0) {
-      return 0;
+  /// Transfer ownership of the values already decoded to the caller.
+  std::shared_ptr<ResizableBuffer> ReleaseValues(MemoryPool* pool) {
+    if (is_void()) {
+      return nullptr;
     }
-    return repetition_level_decoder_.Decode(static_cast<int>(batch_size), 
levels);
+
+    auto result = values_;
+    const auto byte_count = BytesCounter::bytes_for_values(values_count());
+    PARQUET_THROW_NOT_OK(result->Resize(byte_count, /*shrink_to_fit=*/true));
+    values_ = AllocateBuffer(pool);
+    ValueSinkCursor::reset();
+    return result;
   }
 
-  // Advance to the next data page
-  bool ReadNewPage() {
-    // Loop until we find the next data page.
-    while (true) {
-      current_page_ = pager_->NextPage();
-      if (!current_page_) {
-        // EOS
-        return false;
+  /// Exponentially reserve more capacity if needed.
+  void ReserveValues(int64_t extra_values) {
+    const auto old_capacity = fit_capacity_for_extra(extra_values);
+    if (capacity() > old_capacity && !is_void()) {
+      const auto byte_count = BytesCounter::bytes_for_values(capacity());
+      PARQUET_THROW_NOT_OK(values_->Resize(byte_count, 
/*shrink_to_fit=*/false));
+      if constexpr (BytesCounter::kZeroNewCapacity) {
+        const auto old_byte_count = 
BytesCounter::bytes_for_values(old_capacity);
+        std::memset(data() + old_byte_count, 0,
+                    static_cast<std::size_t>(byte_count - old_byte_count));
       }
+    }
+  }
 
-      if (current_page_->type() == PageType::DICTIONARY_PAGE) {
-        ConfigureDictionary(static_cast<const 
DictionaryPage*>(current_page_.get()));
-        continue;
-      } else if (current_page_->type() == PageType::DATA_PAGE) {
-        const auto* page = static_cast<const DataPageV1*>(current_page_.get());
-        const int64_t levels_byte_size = InitializeLevelDecoders(
-            *page, page->repetition_level_encoding(), 
page->definition_level_encoding());
-        InitializeDataDecoder(*page, levels_byte_size);
-        return true;
-      } else if (current_page_->type() == PageType::DATA_PAGE_V2) {
-        const auto* page = static_cast<const DataPageV2*>(current_page_.get());
-        int64_t levels_byte_size = InitializeLevelDecodersV2(*page);
-        InitializeDataDecoder(*page, levels_byte_size);
-        return true;
+  void ResetValues() {
+    if (!is_void()) {
+      PARQUET_THROW_NOT_OK(values_->Resize(0, /*shrink_to_fit=*/false));
+      ValueSinkCursor::reset();
+    }
+  }
+
+  void DebugPrintState() {
+    const T* vals = data();
+    for (int64_t i = 0; i < values_count(); ++i) {
+      if constexpr (can_cout<T>) {
+        std::cout << vals[i] << ' ';
       } else {
-        // We don't know what this page type is. We're allowed to skip non-data
-        // pages.
-        continue;
+        std::cout << "? ";
       }
     }
-    return true;
   }
 
-  void ConfigureDictionary(const DictionaryPage* page) {
-    int encoding = static_cast<int>(page->encoding());
-    if (page->encoding() == Encoding::PLAIN_DICTIONARY ||
-        page->encoding() == Encoding::PLAIN) {
-      encoding = static_cast<int>(Encoding::RLE_DICTIONARY);
+ protected:
+  /// Decode values contiguously into the sink buffer.
+  template <typename Int>
+  Int ReadFromCallback(auto&& decode, Int batch_size) {
+    if (is_void()) {
+      return {};
     }
 
-    auto it = decoders_.find(encoding);
-    if (it != decoders_.end()) {
-      throw ParquetException("Column cannot have more than one dictionary.");
-    }
+    ReserveValues(batch_size);
+    const Int decoded = decode();
+    set_values_count(values_count() + decoded);
+    return decoded;
+  }
 
-    if (page->encoding() == Encoding::PLAIN_DICTIONARY ||
-        page->encoding() == Encoding::PLAIN) {
-      auto dictionary = MakeTypedDecoder<DType>(Encoding::PLAIN, descr_, 
pool_);
-      dictionary->SetData(page->num_values(), page->data(), page->size());
-
-      // The dictionary is fully decoded during DictionaryDecoder::Init, so the
-      // DictionaryPage buffer is no longer required after this step
-      //
-      // TODO(wesm): investigate whether this all-or-nothing decoding of the
-      // dictionary makes sense and whether performance can be improved
-
-      std::unique_ptr<DictDecoder<DType>> decoder = 
MakeDictDecoder<DType>(descr_, pool_);
-      decoder->SetDict(dictionary.get());
-      decoders_[encoding] =
-          
std::unique_ptr<DecoderType>(dynamic_cast<DecoderType*>(decoder.release()));
-    } else {
-      ParquetException::NYI("only plain dictionary encoding has been 
implemented");
-    }
+  value_type* write_start() { return data() + values_count(); }
 
-    new_dictionary_ = true;
-    current_decoder_.SetDecoder(decoders_[encoding].get());
-    ARROW_DCHECK(current_decoder_);
-  }
-
-  // Initialize repetition and definition level decoders on the next data page.
-
-  // If the data page includes repetition and definition levels, we
-  // initialize the level decoders and return the number of encoded level 
bytes.
-  // The return value helps determine the number of bytes in the encoded data.
-  int64_t InitializeLevelDecoders(const DataPage& page,
-                                  Encoding::type repetition_level_encoding,
-                                  Encoding::type definition_level_encoding) {
-    // Read a data page.
-    num_buffered_values_ = page.num_values();
-
-    // Have not decoded any values from the data page yet
-    num_decoded_values_ = 0;
-
-    const uint8_t* buffer = page.data();
-    int32_t levels_byte_size = 0;
-    int32_t max_size = page.size();
-
-    // Data page Layout: Repetition Levels - Definition Levels - encoded 
values.
-    // Levels are encoded as rle or bit-packed.
-    // Init repetition levels
-    if (max_rep_level() > 0) {
-      int32_t rep_levels_bytes = repetition_level_decoder_.SetData(
-          repetition_level_encoding, max_rep_level(),
-          static_cast<int>(num_buffered_values_), buffer, max_size);
-      buffer += rep_levels_bytes;
-      levels_byte_size += rep_levels_bytes;
-      max_size -= rep_levels_bytes;
-    }
-    // TODO figure a way to set max_def_level_ to 0
-    // if the initial value is invalid
-
-    // Init definition levels
-    if (max_def_level() > 0) {
-      int32_t def_levels_bytes = definition_level_decoder_.SetData(
-          definition_level_encoding, max_def_level(),
-          static_cast<int>(num_buffered_values_), buffer, max_size);
-      levels_byte_size += def_levels_bytes;
-      max_size -= def_levels_bytes;
-    }
+ private:
+  std::shared_ptr<::arrow::ResizableBuffer> values_;
+};
 
-    return levels_byte_size;
+/*********************
+ *  ValueSinkBuffer  *
+ *********************/
+
+/// A data sink for decoded values that writes into an Arrow buffer.
+template <typename T>
+class ValueSinkBuffer : public DataSinkBuffer<T> {
+ public:
+  using Base = DataSinkBuffer<T>;
+  using value_type = T;
+
+  explicit ValueSinkBuffer(MemoryPool* pool) : Base(AllocateBuffer(pool)) {}
+
+  void OnNewDictionary(auto& /* decoder */) {}
+
+  /// Decode values contiguously into the sink buffer.
+  void ReadValuesDense(auto& decoder, int32_t batch_size) {
+    const auto decoded = Base::ReadFromCallback(
+        [&, this]() { return decoder.Decode(Base::write_start(), batch_size); 
},
+        batch_size);
+    CheckNumberDecoded(decoded, batch_size);
+  }
+
+  /// Decode values according to a validity bitmap.
+  ///
+  /// The values are decoded in the buffer where the validity bit is set. The 
values
+  /// associated to unset bits are skipped and can be set to any arbitrary 
value.
+  void ReadValuesSpaced(auto& decoder, int32_t batch_size, int32_t null_count,
+                        const uint8_t* valid_bits, int64_t valid_bits_offset) {
+    const auto decoded = Base::ReadFromCallback(
+        [&, this]() {
+          return decoder.DecodeSpaced(Base::write_start(), batch_size, 
null_count,
+                                      valid_bits, valid_bits_offset);
+        },
+        batch_size);
+    CheckNumberDecoded(decoded, batch_size);
   }
+};
 
-  int64_t InitializeLevelDecodersV2(const DataPageV2& page) {
-    // Read a data page.
-    num_buffered_values_ = page.num_values();
+/*********************
+ *  LevelSinkBuffer  *
+ *********************/
 
-    // Have not decoded any values from the data page yet
-    num_decoded_values_ = 0;
-    const uint8_t* buffer = page.data();
+/// A data sink for decoded levels in int16_t that writes into an Arrow buffer.
+class LevelSinkBuffer : public DataSinkBuffer<int16_t> {
+ public:
+  using Base = DataSinkBuffer<int16_t>;
 
-    const int64_t total_levels_length =
-        static_cast<int64_t>(page.repetition_levels_byte_length()) +
-        page.definition_levels_byte_length();
+  /// Create a noop sink
+  static LevelSinkBuffer MakeVoid() { return LevelSinkBuffer(nullptr); }
 
-    if (total_levels_length > page.size()) {
-      throw ParquetException("Data page too small for levels (corrupt 
header?)");
-    }
+  /// Create a stateful sink with an allocated buffer
+  static LevelSinkBuffer MakeAllocated(MemoryPool* pool) {
+    return LevelSinkBuffer(AllocateBuffer(pool));
+  }
 
-    if (max_rep_level() > 0) {
-      repetition_level_decoder_.SetDataV2(page.repetition_levels_byte_length(),
-                                          max_rep_level(),
-                                          
static_cast<int>(num_buffered_values_), buffer);
-    }
-    // ARROW-17453: Even if max_rep_level_ is 0, there may still be
-    // repetition level bytes written and/or reported in the header by
-    // some writers (e.g. Athena)
-    buffer += page.repetition_levels_byte_length();
-
-    if (max_def_level() > 0) {
-      definition_level_decoder_.SetDataV2(page.definition_levels_byte_length(),
-                                          max_def_level(),
-                                          
static_cast<int>(num_buffered_values_), buffer);
-    }
+  /// Decode levels into the sink buffer.
+  void Decode(auto& decoder, int32_t batch_size) {
+    const auto decoded = Base::ReadFromCallback(
+        [&, this]() { return decoder.Decode(Base::write_start(), batch_size); 
},
+        batch_size);
+    CheckValidLevelCount(decoded == batch_size);
+  }
+
+ private:
+  explicit LevelSinkBuffer(auto buffer) : Base(std::move(buffer)) {}
+};
+
+/************************
+ *  ValiditySinkBuffer  *
+ ************************/
+
+struct BytesCounterForBits {
+  /// Appending bits via read-modify-write may read uniitialized memory.
+  /// Must zero new capacity to avoid Valgrind/MSAN warnings.
+  static constexpr bool kZeroNewCapacity = true;
 
-    return total_levels_length;
+  static int64_t bytes_for_values(int64_t num_values) {
+    return bit_util::BytesForBits(num_values);
   }
+};
 
-  // Get a decoder object for this page or create a new decoder if this is the
-  // first page with this encoding.
-  void InitializeDataDecoder(const DataPage& page, int64_t levels_byte_size) {
-    const uint8_t* buffer = page.data() + levels_byte_size;
-    const int64_t data_size = page.size() - levels_byte_size;
+/// A data sink for a validity bitmap that writes into an Arrow buffer.
+class ValiditySinkBuffer : private DataSinkBuffer<uint8_t, 
BytesCounterForBits> {
+ public:
+  using Base = DataSinkBuffer<uint8_t, BytesCounterForBits>;
 
-    if (data_size < 0) {
-      throw ParquetException("Page smaller than size of encoded levels");
-    }
+  using Base::data;
+  using Base::is_void;
+  using Base::ReleaseValues;
+  using Base::ReserveValues;
+  using Base::ResetValues;
+  using Base::values_count;
 
-    Encoding::type encoding = page.encoding();
-    if (IsDictionaryIndexEncoding(encoding)) {
-      // Normalizing the PLAIN_DICTIONARY to RLE_DICTIONARY encoding
-      // in decoder.
-      encoding = Encoding::RLE_DICTIONARY;
-    }
+  /// Create a noop sink
+  static ValiditySinkBuffer MakeVoid() { return ValiditySinkBuffer(nullptr); }
 
-    auto it = decoders_.find(static_cast<int>(encoding));
-    if (it != decoders_.end()) {
-      ARROW_DCHECK(it->second.get() != nullptr);
-      current_decoder_.SetDecoder(it->second.get());
-    } else {
-      switch (encoding) {
-        case Encoding::PLAIN:
-        case Encoding::BYTE_STREAM_SPLIT:
-        case Encoding::RLE:
-        case Encoding::DELTA_BINARY_PACKED:
-        case Encoding::DELTA_BYTE_ARRAY:
-        case Encoding::DELTA_LENGTH_BYTE_ARRAY: {
-          auto decoder = MakeTypedDecoder<DType>(encoding, descr_, pool_);
-          current_decoder_.SetDecoder(decoder.get());
-          decoders_[static_cast<int>(encoding)] = std::move(decoder);
-          break;
-        }
+  /// Create a stateful sink with an allocated buffer
+  static ValiditySinkBuffer MakeAllocated(MemoryPool* pool) {
+    return ValiditySinkBuffer(AllocateBuffer(pool));
+  }
 
-        case Encoding::RLE_DICTIONARY:
-          throw ParquetException("Dictionary page must be before data page.");
+  struct ReadResult {
+    int64_t values_read = 0;
+    int64_t null_count = 0;
+  };
 
-        default:
-          throw ParquetException("Unknown encoding type.");
-      }
+  /// Write into the bitmap from already decoded levels.
+  ///
+  /// This method is used when additional computation using the levels is 
needed.
+  /// It re-encodes definition levels as a validity bitmap.
+  ///
+  /// Unlike the other sink methods, the number of values written is not known 
in
+  /// advance and is returned to the caller: for repeated or nested columns 
some
+  /// definition levels describe empty or absent lists that produce no leaf 
value.
+  ReadResult ReadFromDefLevels(const int16_t* def_levels, int64_t 
num_def_levels,
+                               const internal::LevelInfo& level_info) {
+    ReadResult out{};
+
+    Base::ReadFromCallback(
+        [&, this]() {
+          internal::ValidityBitmapInputOutput validity_io{};
+          validity_io.values_read_upper_bound = num_def_levels;
+          validity_io.valid_bits = data();
+          validity_io.valid_bits_offset = values_count();
+          DefLevelsToBitmap(def_levels, num_def_levels, level_info, 
&validity_io);
+          ARROW_DCHECK_GE(validity_io.values_read, 0);
+          ARROW_DCHECK_GE(validity_io.null_count, 0);
+
+          out.values_read = validity_io.values_read;
+          out.null_count = validity_io.null_count;
+
+          // Advance by the number of leaf values (one validity bit each), not 
by the
+          // number of definition levels: for repeated/nested columns some def 
levels
+          // describe empty or absent lists that produce no leaf value, so
+          // values_read <= num_def_levels.
+          return validity_io.values_read;
+        },
+        num_def_levels);
+
+    return out;
+  }
+
+  /// Write into the bitmap from a bitmap-compatible decoder.
+  ///
+  /// This is a special optimization when the max definition level is one, in 
which case
+  /// definition levels are already a validity bitmap.
+  ///
+  /// @return the number of null values written.
+  int32_t ReadFromDecoder(PageLevelToBitmapDecoder& decoder, int32_t 
batch_size) {
+    int32_t null_count = 0;
+
+    const auto decoded = Base::ReadFromCallback(
+        [&, this]() {
+          using ::arrow::util::BitmapSpanMut;
+
+          const auto write_pos = BitmapSpanMut(
+              /* data= */ data() + values_count() / 8,
+              /* bit_offset= */ static_cast<int32_t>(values_count() % 8));
+          const auto decoded = decoder.Decode(write_pos, batch_size);
+          const auto non_null = ::arrow::internal::CountSetBits(
+              write_pos.data(), write_pos.bit_start(), decoded);
+          ARROW_DCHECK_LE(non_null, decoded);
+
+          null_count = static_cast<int32_t>(decoded - non_null);
+          return decoded;
+        },
+        batch_size);
+    CheckValidLevelCount(decoded == batch_size);
+
+    return null_count;
+  }
+
+ private:
+  explicit ValiditySinkBuffer(auto buffer) : Base(std::move(buffer)) {}
+};
+
+/***********************
+ *  ColumnChunkReader  *
+ ***********************/
+
+/// Initialize repetition and definition level decoders on the given data page.
+///
+/// If the data page includes repetition and definition levels, we initialize 
the level
+/// decoders and return the number of encoded level bytes.
+/// The return value helps determine the number of bytes in the encoded data.
+int64_t InitializeV1Levels(const DataPageV1& page, auto& def_dec, auto& 
rep_dec) {
+  const auto num_values = static_cast<int>(page.num_values());
+
+  const uint8_t* buffer = page.data();
+  int32_t levels_byte_size = 0;
+  int32_t max_size = page.size();
+
+  if (const auto max_rep_lvl = rep_dec.max_level(); max_rep_lvl > 0) {
+    const int32_t rep_levels_bytes =
+        rep_dec.SetDataV1(page.repetition_level_encoding(), buffer, max_size,
+                          {.value_count = num_values, .max_level = 
max_rep_lvl});
+    buffer += rep_levels_bytes;
+    levels_byte_size += rep_levels_bytes;
+    max_size -= rep_levels_bytes;
+  }
+
+  if (const auto max_def_lvl = def_dec.max_level(); max_def_lvl > 0) {
+    const int32_t def_levels_bytes =
+        def_dec.SetDataV1(page.definition_level_encoding(), buffer, max_size,
+                          {.value_count = num_values, .max_level = 
max_def_lvl});
+    levels_byte_size += def_levels_bytes;
+    max_size -= def_levels_bytes;
+  }
+
+  return levels_byte_size;
+}
+
+/// Initialize repetition and definition level decoders on the given data page.
+///
+/// If the data page includes repetition and definition levels, we initialize 
the level
+/// decoders and return the number of encoded level bytes.
+/// The return value helps determine the number of bytes in the encoded data.
+int64_t InitializeV2Levels(const DataPageV2& page, auto& def_dec, auto& 
rep_dec) {
+  const auto num_values = static_cast<int>(page.num_values());
+
+  const int64_t total_levels_length =
+      static_cast<int64_t>(page.repetition_levels_byte_length()) +
+      page.definition_levels_byte_length();
+  if (total_levels_length > page.size()) {
+    throw ParquetException("Data page too small for levels (corrupt header?)");
+  }
+
+  const uint8_t* buffer = page.data();
+
+  if (const auto max_rep_lvl = rep_dec.max_level(); max_rep_lvl > 0) {
+    rep_dec.SetDataV2(buffer, page.repetition_levels_byte_length(),
+                      {.value_count = num_values, .max_level = max_rep_lvl});
+  }
+  // ARROW-17453: Even if max_rep_level_ is 0, there may still be
+  // repetition level bytes written and/or reported in the header by
+  // some writers (e.g. Athena)
+  buffer += page.repetition_levels_byte_length();
+
+  if (const auto max_def_lvl = def_dec.max_level(); max_def_lvl > 0) {
+    def_dec.SetDataV2(buffer, page.definition_levels_byte_length(),
+                      {.value_count = num_values, .max_level = max_def_lvl});
+  }
+
+  return total_levels_length;
+}
+
+/// A base reader for iterating through the pages of a column chunk.
+///
+/// It sets the decoders and level decoders appropriately while also providing 
basic
+/// level counting. The number of records depends on the schema: an optional or
+/// repeated field will use more levels to encode a single record into zero, 
one, or
+/// more "values" from its encoder.
+template <typename Traits>
+class ColumnChunkReader {
+ public:
+  using DType = typename Traits::DType;
+  using DefLevelDecoder = typename Traits::DefLevelDecoder;
+  using RepLevelDecoder = typename Traits::RepLevelDecoder;
+  using value_type = typename DType::c_type;
+
+  ColumnChunkReader(const ColumnDescriptor* descr, ::arrow::MemoryPool* pool,
+                    DefLevelDecoder def_levels_decoder,
+                    RepLevelDecoder rep_levels_decoder)
+      : current_decoder_(pool),
+        def_levels_decoder_(std::move(def_levels_decoder)),
+        rep_levels_decoder_(std::move(rep_levels_decoder)),
+        descr_(descr),
+        pool_(pool) {}
+
+  void SetPageReader(std::unique_ptr<PageReader> reader) {
+    current_page_ = nullptr;
+    pager_ = std::move(reader);
+    current_decoder_.SetDecoder(nullptr);
+    decoders_.clear();
+    current_encoding_ = Encoding::UNKNOWN;
+  }
+
+  bool HasPageReader() const { return pager_ != nullptr; }
+
+  /// Return true if there is more data.
+  ///
+  /// If the current page is exhausted, it will process more pages until some 
data
+  /// page is found.
+  bool EnsureDataPage();
+
+  /// Check the encoding of the current page or throw an exception.
+  void CheckEncodingIs(Encoding::type encoding);
+
+  /// Read the current dictionary.
+  ///
+  /// Tries to advance to the initial dictionary page if a the beginning of a 
column
+  /// chunk. Throws if the column chunk is not dictionary encoded.
+  const value_type* ReadDictionary(int32_t* dictionary_length);
+
+  int32_t ReadDefinitionLevels(int32_t batch_size, int16_t* levels) {
+    if (max_def_level() == 0) {
+      return 0;
     }
-    current_encoding_ = encoding;
-    current_decoder_->SetData(static_cast<int>(num_buffered_values_), buffer,
-                              static_cast<int>(data_size));
+    return def_levels_decoder_.Decode(levels, batch_size);
   }
 
-  // Available values in the current data page, value includes repeated values
-  // and nulls.
-  int64_t available_values_current_page() const {
-    return num_buffered_values_ - num_decoded_values_;
+  int32_t ReadRepetitionLevels(int32_t batch_size, int16_t* levels) {
+    if (max_rep_level() == 0) {
+      return 0;
+    }
+    return rep_levels_decoder_.Decode(levels, batch_size);
   }
 
-  int16_t max_def_level() const {
-    // max level indirectly part of this object storage
-    return definition_level_decoder_.max_level();
+  // Available values in the current data page, value includes repeated values 
and nulls.
+  int32_t available_values_current_page() const {
+    const int32_t out = num_buffered_values_ - num_decoded_values_;
+    ARROW_DCHECK_GE(out, 0);
+    return out;
   }
 
-  int16_t max_rep_level() const {
-    // max level indirectly part of this object storage
-    return repetition_level_decoder_.max_level();
+  /// Return the current decoder as a dict decoder if page is dictionary 
encoded
+  DictDecoder<DType>* current_dict_decoder() {
+    return dynamic_cast<DictDecoder<DType>*>(this->current_decoder_.get());
   }
 
-  const ColumnDescriptor* descr_;
+  int16_t max_def_level() const { return def_levels_decoder_.max_level(); }
 
-  std::unique_ptr<PageReader> pager_;
-  std::shared_ptr<Page> current_page_;
+  int16_t max_rep_level() const { return rep_levels_decoder_.max_level(); }
+
+  void MarkValuesAsConsumed(int32_t num_values) { num_decoded_values_ += 
num_values; }
 
-  // No data set if full schema for this field has no optional or repeated 
elements
-  LevelDecoder definition_level_decoder_;
+  int64_t Skip(int64_t num_values_to_skip);
+
+ protected:
+  SkippableTypedDecoder<DType> current_decoder_;
+  DefLevelDecoder def_levels_decoder_;
+  RepLevelDecoder rep_levels_decoder_;
+  const ColumnDescriptor* descr_;
+  ::arrow::MemoryPool* pool_;
 
-  // No data set for flat schemas.
-  LevelDecoder repetition_level_decoder_;
+ private:
+  using DecoderType = TypedDecoder<DType>;
 
+  // Map of encoding type to the respective decoder object. For example, a
+  // column chunk's data pages may include both dictionary-encoded and
+  // plain-encoded data.
+  std::unordered_map<int, std::unique_ptr<DecoderType>> decoders_;
+  std::unique_ptr<PageReader> pager_;
+  std::shared_ptr<DataPage> current_page_;
   // The total number of values stored in the data page. This is the maximum of
   // the number of encoded definition levels or encoded values. For
   // non-repeated, required columns, this is equal to the number of encoded
   // values. For repeated or optional values, there may be fewer data values
   // than levels, and this tells you how many encoded levels there are in that
   // case.
-  int64_t num_buffered_values_ = 0;
-
+  int32_t num_buffered_values_ = 0;
   // The number of values from the current data page that have been decoded
   // into memory or skipped over.
-  int64_t num_decoded_values_ = 0;
-
-  ::arrow::MemoryPool* pool_;
-
-  using DecoderType = TypedDecoder<DType>;
-  SkippableTypedDecoder<DType, kSkipScratchBatchSize> current_decoder_;
+  int32_t num_decoded_values_ = 0;
   Encoding::type current_encoding_ = Encoding::UNKNOWN;
 
-  /// Flag to signal when a new dictionary has been set, for the benefit of
-  /// DictionaryRecordReader
-  bool new_dictionary_ = false;
+  // Advance to the next data page
+  bool ReadNewPage();
 
-  // The exposed encoding
-  ExposedEncoding exposed_encoding_ = ExposedEncoding::NO_ENCODING;
+  void ConfigureDictionary(const DictionaryPage* page);
 
-  // Map of encoding type to the respective decoder object. For example, a
-  // column chunk's data pages may include both dictionary-encoded and
-  // plain-encoded data.
-  std::unordered_map<int, std::unique_ptr<DecoderType>> decoders_;
+  template <typename DP>
+  void InitializeDataPage(std::shared_ptr<DP> page);
+
+  // Get a decoder object for this page or create a new decoder if this is the
+  // first page with this encoding.
+  void InitializeDataDecoder(const DataPage& page, int64_t levels_byte_size);
 
-  void ConsumeBufferedValues(int64_t num_values) { num_decoded_values_ += 
num_values; }
+  /// Advance both level decoders by num_levels, without materializing them.
+  ///
+  /// @return The number of values present (non-null) among them, which is the
+  ///         number of values to skip in the data decoder.
+  int32_t AdvanceLevels(int32_t num_levels);
 };
 
-// ----------------------------------------------------------------------
-// TypedColumnReader implementations
+/**************************************
+ *  ColumnChunkReader Implementation  *
+ **************************************/
+
+template <typename Traits>
+void ColumnChunkReader<Traits>::CheckEncodingIs(Encoding::type encoding) {
+  if (current_encoding_ != encoding) {
+    throw ParquetException("Unexpected data page encoding. Expected ",
+                           EncodingToString(encoding), ", got ",
+                           EncodingToString(current_encoding_));
+  }
+}
+
+template <typename Traits>
+auto ColumnChunkReader<Traits>::ReadDictionary(int32_t* dictionary_length)
+    -> const value_type* {
+  if (!this->current_decoder_ && !this->EnsureDataPage()) {
+    *dictionary_length = 0;
+    return nullptr;
+  }
+  // Verify the current data page is dictionary encoded.
+  this->CheckEncodingIs(Encoding::RLE_DICTIONARY);
+  const value_type* dictionary = nullptr;
+  current_dict_decoder()->GetDictionary(&dictionary, dictionary_length);
+  return dictionary;
+}
+
+template <typename Traits>
+bool ColumnChunkReader<Traits>::EnsureDataPage() {
+  // 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;
+    }
+  }
+  return true;
+}
+
+template <typename Traits>
+bool ColumnChunkReader<Traits>::ReadNewPage() {
+  // Loop until we find the next data page.
+  while (true) {
+    std::shared_ptr<Page> page = pager_->NextPage();
+    if (!page) {
+      // EOS
+      return false;
+    }
+
+    if (page->type() == PageType::DICTIONARY_PAGE) {
+      ConfigureDictionary(static_cast<const DictionaryPage*>(page.get()));
+      continue;
+    } else if (page->type() == PageType::DATA_PAGE) {
+      InitializeDataPage(std::static_pointer_cast<DataPageV1>(page));
+      return true;
+    } else if (page->type() == PageType::DATA_PAGE_V2) {
+      InitializeDataPage(std::static_pointer_cast<DataPageV2>(page));
+      return true;
+    } else {
+      // We don't know what this page type is. We're allowed to skip non-data
+      // pages.
+      continue;
+    }
+  }
+  return true;
+}
+
+template <typename Traits>
+void ColumnChunkReader<Traits>::ConfigureDictionary(const DictionaryPage* 
page) {
+  int encoding = static_cast<int>(page->encoding());
+  if (page->encoding() == Encoding::PLAIN_DICTIONARY ||
+      page->encoding() == Encoding::PLAIN) {
+    encoding = static_cast<int>(Encoding::RLE_DICTIONARY);
+  }
+
+  auto it = decoders_.find(encoding);
+  if (it != decoders_.end()) {
+    throw ParquetException("Column cannot have more than one dictionary.");
+  }
+
+  if (page->encoding() == Encoding::PLAIN_DICTIONARY ||
+      page->encoding() == Encoding::PLAIN) {
+    auto dictionary = MakeTypedDecoder<DType>(Encoding::PLAIN, descr_, pool_);
+    dictionary->SetData(page->num_values(), page->data(), page->size());
+
+    // The dictionary is fully decoded during DictionaryDecoder::Init, so the
+    // DictionaryPage buffer is no longer required after this step
+    //
+    // TODO(wesm): investigate whether this all-or-nothing decoding of the
+    // dictionary makes sense and whether performance can be improved
+
+    std::unique_ptr<DictDecoder<DType>> decoder = 
MakeDictDecoder<DType>(descr_, pool_);
+    decoder->SetDict(dictionary.get());
+    decoders_[encoding] = std::move(decoder);
+  } else {
+    ParquetException::NYI("only plain dictionary encoding has been 
implemented");
+  }
+
+  current_decoder_.SetDecoder(decoders_[encoding].get());
+  ARROW_DCHECK(current_decoder_);
+}
+
+template <typename Traits>
+template <typename DP>
+void ColumnChunkReader<Traits>::InitializeDataPage(std::shared_ptr<DP> page) {
+  ARROW_DCHECK_NE(page, nullptr);
+  current_page_ = std::static_pointer_cast<DataPage>(page);
+  int64_t byte_size = 0;
+  if constexpr (std::is_same_v<DP, DataPageV1>) {
+    byte_size = InitializeV1Levels(*page, def_levels_decoder_, 
rep_levels_decoder_);
+  } else if constexpr (std::is_same_v<DP, DataPageV2>) {
+    byte_size = InitializeV2Levels(*page, def_levels_decoder_, 
rep_levels_decoder_);
+  } else {
+    // Some compiler do not support `static_assert(false, ...)` in a discarded 
branch
+    static_assert(!std::is_same_v<DP, DP>, "Unknown data page type");
+  }
+  num_buffered_values_ = page->num_values();
+  num_decoded_values_ = 0;
+  InitializeDataDecoder(*page, byte_size);
+}
+
+template <typename Traits>
+void ColumnChunkReader<Traits>::InitializeDataDecoder(const DataPage& page,
+                                                      int64_t 
levels_byte_size) {
+  const uint8_t* buffer = page.data() + levels_byte_size;
+  const int64_t data_size = page.size() - levels_byte_size;
+
+  if (data_size < 0) {
+    throw ParquetException("Page smaller than size of encoded levels");
+  }
+
+  Encoding::type encoding = page.encoding();
+  if (IsDictionaryIndexEncoding(encoding)) {
+    // Normalizing the PLAIN_DICTIONARY to RLE_DICTIONARY encoding
+    // in decoder.
+    encoding = Encoding::RLE_DICTIONARY;
+  }
+
+  auto it = decoders_.find(static_cast<int>(encoding));
+  if (it != decoders_.end()) {
+    ARROW_DCHECK(it->second.get() != nullptr);
+    current_decoder_.SetDecoder(it->second.get());
+  } else {
+    switch (encoding) {
+      case Encoding::PLAIN:
+      case Encoding::BYTE_STREAM_SPLIT:
+      case Encoding::RLE:
+      case Encoding::DELTA_BINARY_PACKED:
+      case Encoding::DELTA_BYTE_ARRAY:
+      case Encoding::DELTA_LENGTH_BYTE_ARRAY: {
+        auto decoder = MakeTypedDecoder<DType>(encoding, descr_, pool_);
+        current_decoder_.SetDecoder(decoder.get());
+        decoders_[static_cast<int>(encoding)] = std::move(decoder);
+        break;
+      }
+
+      case Encoding::RLE_DICTIONARY:
+        throw ParquetException("Dictionary page must be before data page.");
+
+      default:
+        throw ParquetException("Unknown encoding type.");
+    }
+  }
+  current_encoding_ = encoding;
+  current_decoder_->SetData(static_cast<int>(num_buffered_values_), buffer,
+                            static_cast<int>(data_size));
+}
+
+template <typename Traits>
+int32_t ColumnChunkReader<Traits>::AdvanceLevels(int32_t num_levels) {
+  int max_count = num_levels;
+  // Advance the definition levels, counting how many correspond to present
+  // (non-null) values that must be skipped in the data decoder.
+  if (this->max_def_level() > 0) {
+    const auto count = def_levels_decoder_.CountUpTo(this->max_def_level(), 
num_levels);

Review Comment:
   `def_levels_decoder_` is now a `PageLevelDecoder`, whose `CountUpTo` accepts 
a `bool`; passing `max_def_level()` here therefore converts every nonzero max 
level to `true` and counts definition level 1, not the actual maximum level. 
For nested/optional columns with `max_definition_level() > 1`, `SkipRecords` 
will skip the wrong number of physical values and desynchronize the value 
decoder. Keep the generic level-counting path level-valued (or add an 
appropriate overload for the bitmap specialization) before using it here.



##########
cpp/src/parquet/level_decoder_internal.h:
##########
@@ -0,0 +1,314 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#pragma once
+
+#include <cstdint>
+#include <sstream>
+#include <variant>
+
+#include "arrow/util/bit_util.h"
+#include "arrow/util/int_util_overflow.h"
+#include "arrow/util/logging.h"
+#include "arrow/util/macros.h"
+#include "arrow/util/rle_bitmap_internal.h"
+#include "arrow/util/rle_encoding_internal.h"
+#include "arrow/util/ubsan.h"
+#include "parquet/exception.h"
+#include "parquet/level_comparison.h"
+#include "parquet/types.h"
+
+namespace parquet {
+
+/*******************************
+ *  PageLevelDecoder Adapters  *
+ *******************************/
+
+/// BitPackedDecoder adapter for PageLevelDecoder.
+struct LevelBitPackedDecoder : ::arrow::util::BitPackedDecoder<int16_t> {
+  using Base = BitPackedDecoder<int16_t>;
+
+  LevelBitPackedDecoder(const uint8_t* data, int32_t data_size, const auto& 
params)
+      : Base(data, data_size,
+             /* value_bit_width= */ ::arrow::bit_util::Log2(params.max_level + 
1),
+             params.value_count) {}
+
+  int32_t GetBatch(int16_t* out, int32_t batch_size, int16_t max_level);
+};
+
+/// RleBitPackedDecoder adapter for PageLevelDecoder.
+struct LevelRleBitPackedDecoder : ::arrow::util::RleBitPackedDecoder<int16_t> {
+  using Base = RleBitPackedDecoder<int16_t>;
+
+  LevelRleBitPackedDecoder(const uint8_t* data, int32_t data_size, const auto& 
params)
+      : Base(data, data_size,
+             /* value_bit_width= */ ::arrow::bit_util::Log2(params.max_level + 
1)) {}
+
+  int32_t GetBatch(int16_t* out, int32_t batch_size, int16_t max_level);
+};
+
+/**************************************
+ *  PageLevelDecoder Bitmap Adapters  *
+ **************************************/
+
+struct LevelToBitmapBitPackedDecoder : ::arrow::util::BitPackedToBitmapDecoder 
{
+  using Base = BitPackedToBitmapDecoder;
+
+  LevelToBitmapBitPackedDecoder(const uint8_t* data, int32_t data_size,
+                                const auto& params)
+      : Base(data, data_size, params.value_count) {
+    ARROW_DCHECK_EQ(params.max_level, 1);
+  }
+
+  template <typename Out>
+  int32_t GetBatch(Out&& out, int32_t batch_size, int16_t /* max_level */) {
+    return Base::GetBatch(out, batch_size);
+  }
+};
+
+struct LevelToBitmapRleBitPackedDecoder : 
::arrow::util::RleBitPackedToBitmapDecoder {
+  using Base = RleBitPackedToBitmapDecoder;
+
+  LevelToBitmapRleBitPackedDecoder(const uint8_t* data, int32_t data_size,
+                                   const auto& params)
+      : Base(data, data_size) {
+    ARROW_DCHECK_EQ(params.max_level, 1);
+  }
+
+  template <typename Out>
+  int32_t GetBatch(Out&& out, int32_t batch_size, int16_t /* max_level */) {
+    return Base::GetBatch(out, batch_size);
+  }
+};
+
+/**********************
+ *  PageLevelDecoder  *
+ **********************/
+
+/// Decoder for repetition or definition level.
+///
+/// This decoder is used with either a deprecated bit packed (`BIT_PACKED = 4`)
+/// encoding or a mixed bit packed and RLE one (`RLE = 3`).
+/// Because it takes as input a single buffer, `SetData` and `Decode` are 
typically
+/// used on each of the Parquet `DataPage`.
+/// The number of levels is guaranteed to fit into an `int32_t` by the 
specification.
+///
+/// @see https://research.google.com/pubs/archive/36632.pdf
+template <typename BitDecoder = LevelBitPackedDecoder,
+          typename RleDecoder = LevelRleBitPackedDecoder>
+class PageLevelDecoder {
+ public:
+  struct DataParams {
+    int32_t value_count = 0;
+    int16_t max_level = 0;
+  };
+
+  explicit PageLevelDecoder(int16_t max_level = 0)
+      : decoder_(BitDecoder(nullptr, 0, DataParams{.max_level = max_level})),
+        max_level_(max_level) {}
+
+  /// Initialize the decoder state with new data from a legacy (V1) page.
+  ///
+  /// @return the number of bytes consumed
+  int32_t SetDataV1(Encoding::type encoding, const uint8_t* data, int32_t 
max_data_size,
+                    const DataParams& params);
+
+  /// Initialize the decoder state with new data from a V2 page.
+  ///
+  /// Repetition and definition levels in V2 pages are always RLE encoded.
+  void SetDataV2(const uint8_t* data, int32_t data_size, const DataParams& 
params);
+
+  /// Decode a batch of levels into `out` and return the number of levels 
decoded.
+  template <typename Out>
+  int32_t Decode(Out&& out, int32_t batch_size);
+
+  /// Advance the decoder and throw away decoded levels.
+  int32_t Skip(int32_t batch_size);
+
+  struct CountUpToResult {
+    int32_t matching_count;
+    int32_t processed_count;
+  };
+
+  /// Advance and count the number of occurrences of `value`.
+  ///
+  /// The count is limited to at most the next `batch_size` items.
+  /// @return The matching value count and number of elements that were 
processed.
+  CountUpToResult CountUpTo(bool value, int32_t batch_size);

Review Comment:
   `PageLevelDecoder` is also used for ordinary multi-level definition 
decoders, but this parameter is `bool`. `ColumnChunkReader::AdvanceLevels` 
passes `max_def_level()` here, so a max definition level of 2 or greater is 
converted to `true` and the decoder counts level 1 instead of the actual max 
level; skipping then consumes the wrong number of physical values. Keep this 
API value-typed (for example `int16_t`) and update the out-of-line definition; 
the bitmap specialization can still receive 0/1.



##########
cpp/src/parquet/column_reader.cc:
##########
@@ -1314,969 +1686,1388 @@ namespace internal {
 
 namespace {
 
-template <typename DType>
-class TypedRecordReader : public TypedColumnReaderImpl<DType>,
+/***********************
+ *  TypedRecordReader  *
+ ***********************/
+
+template <typename D>
+struct TypedRecordReaderTraits {
+  using DType = D;
+  using DefLevelDecoder = NewLevelDecoder;
+  using RepLevelDecoder = NewLevelDecoder;
+};
+
+/// General record reader for a given data type.
+///
+/// This historical class can read all repetition, at the cost of increased 
complexity.
+/// The main difficulties are that repeated values will span multiple values 
(sometimes
+/// across data pages) in Parquet, while nulls are not written.
+/// Decoding levels is therefore critical to reconstruct the delimitation 
across records,
+/// but makes optimizing simple cases harder.
+template <typename DType, typename ValueSink, bool kReadDictionary>
+class TypedRecordReader : public 
ColumnChunkReader<TypedRecordReaderTraits<DType>>,
                           virtual public RecordReader {
  public:
   using T = typename DType::c_type;
-  using BASE = TypedColumnReaderImpl<DType>;
-  TypedRecordReader(const ColumnDescriptor* descr, LevelInfo leaf_info, 
MemoryPool* pool,
-                    bool read_dense_for_nullable)
-      // Pager must be set using SetPageReader.
-      : BASE(descr, /* pager = */ nullptr, pool) {
-    leaf_info_ = leaf_info;
-    nullable_values_ = leaf_info_.HasNullableValues();
-    at_record_start_ = true;
-    values_written_ = 0;
-    null_count_ = 0;
-    values_capacity_ = 0;
-    levels_written_ = 0;
-    levels_position_ = 0;
-    levels_capacity_ = 0;
-    read_dense_for_nullable_ = read_dense_for_nullable;
-    // FIXED_LEN_BYTE_ARRAY and BYTE_ARRAY values are not stored in the 
`values_` buffer,
-    // they are read directly as Arrow.
-    uses_values_ = (descr->physical_type() != Type::BYTE_ARRAY &&
-                    descr->physical_type() != Type::FIXED_LEN_BYTE_ARRAY);
-
-    if (uses_values_) {
-      values_ = AllocateBuffer(pool);
-    }
-    valid_bits_ = AllocateBuffer(pool);
-    def_levels_ = AllocateBuffer(pool);
-    rep_levels_ = AllocateBuffer(pool);
-    TypedRecordReader::Reset();
-  }
+  using Base = ColumnChunkReader<TypedRecordReaderTraits<DType>>;
 
-  // Compute the values capacity in bytes for the given number of elements
-  int64_t bytes_for_values(int64_t nitems) const {
-    int64_t type_size = GetTypeByteSize(this->descr_->physical_type());
-    int64_t bytes_for_values = -1;
-    if (MultiplyWithOverflow(nitems, type_size, &bytes_for_values)) {
-      throw ParquetException("Total size of items too large");
+  TypedRecordReader(const ColumnDescriptor* descr, LevelInfo leaf_info, 
MemoryPool* pool,
+                    bool read_dense_for_nullable, ValueSink value_sink)
+      : Base(descr, pool, NewLevelDecoder(descr->max_definition_level()),
+             NewLevelDecoder(descr->max_repetition_level())),
+        value_sink_(std::move(value_sink)),
+        leaf_info_(leaf_info) {
+    if (!read_dense_for_nullable && nullable_values()) {
+      valid_bits_ = ValiditySinkBuffer::MakeAllocated(this->pool_);
     }
-    return bytes_for_values;
-  }
-
-  const void* ReadDictionary(int32_t* dictionary_length) override {
-    if (!this->current_decoder_ && !this->HasNextInternal()) {
-      *dictionary_length = 0;
-      return nullptr;
+    if (this->max_def_level() > 0) {
+      def_levels_ = LevelSinkBuffer::MakeAllocated(pool);
     }
-    // Verify the current data page is dictionary encoded. The 
current_encoding_ should
-    // have been set as RLE_DICTIONARY if the page encoding is RLE_DICTIONARY 
or
-    // PLAIN_DICTIONARY.
-    if (this->current_encoding_ != Encoding::RLE_DICTIONARY) {
-      std::stringstream ss;
-      ss << "Data page is not dictionary encoded. Encoding: "
-         << EncodingToString(this->current_encoding_);
-      throw ParquetException(ss.str());
+    if (this->max_rep_level() > 0) {
+      rep_levels_ = LevelSinkBuffer::MakeAllocated(pool);
     }
-    auto decoder = 
dynamic_cast<DictDecoder<DType>*>(this->current_decoder_.get());
-    const T* dictionary = nullptr;
-    decoder->GetDictionary(&dictionary, dictionary_length);
-    return reinterpret_cast<const void*>(dictionary);
   }
 
-  int64_t ReadRecords(int64_t num_records) override {
-    if (num_records == 0) return 0;
-    // Delimit records, then read values at the end
-    int64_t records_read = 0;
+  int16_t* def_levels() const final { return def_levels_.data(); }
 
-    if (has_values_to_process()) {
-      records_read += ReadRecordData(num_records);
-    }
+  int16_t* rep_levels() const final { return rep_levels_.data(); }
 
-    int64_t level_batch_size = std::max<int64_t>(kMinLevelBatchSize, 
num_records);
-
-    // If we are in the middle of a record, we continue until reaching the
-    // desired number of records or the end of the current record if we've 
found
-    // enough records
-    while (!at_record_start_ || records_read < num_records) {
-      // Is there more data to read in this row group?
-      if (!this->HasNextInternal()) {
-        if (!at_record_start_) {
-          // We ended the row group while inside a record that we haven't seen
-          // the end of yet. So increment the record count for the last record 
in
-          // the row group
-          ++records_read;
-          at_record_start_ = true;
-        }
-        break;
-      }
+  int64_t levels_position() const final { return levels_position_; }
 
-      /// We perform multiple batch reads until we either exhaust the row group
-      /// or observe the desired number of records
-      int64_t batch_size =
-          std::min(level_batch_size, this->available_values_current_page());
+  int64_t levels_written() const final { return def_levels_.values_count(); }
 
-      // No more data in column
-      if (batch_size == 0) {
-        break;
-      }
-
-      if (this->max_def_level() > 0) {
-        ReserveLevels(batch_size);
+  int64_t null_count() const final { return null_count_; }
 
-        int16_t* def_levels = this->def_levels() + levels_written_;
-        int16_t* rep_levels = this->rep_levels() + levels_written_;
+  bool nullable_values() const final { return leaf_info_.HasNullableValues(); }
 
-        if (ARROW_PREDICT_FALSE(this->ReadDefinitionLevels(batch_size, 
def_levels) !=
-                                batch_size)) {
-          throw ParquetException(kErrorRepDefLevelNotMatchesNumValues);
-        }
-        if (this->max_rep_level() > 0) {
-          int64_t rep_levels_read = this->ReadRepetitionLevels(batch_size, 
rep_levels);
-          if (ARROW_PREDICT_FALSE(rep_levels_read != batch_size)) {
-            throw ParquetException(kErrorRepDefLevelNotMatchesNumValues);
-          }
-        }
+  bool read_dictionary() const final { return kReadDictionary; }
 
-        levels_written_ += batch_size;
-        records_read += ReadRecordData(num_records - records_read);
-      } else {
-        // No repetition and definition levels, we can read values directly
-        batch_size = std::min(num_records - records_read, batch_size);
-        records_read += ReadRecordData(batch_size);
-      }
-    }
-
-    return records_read;
+  bool read_dense_for_nullable() const final {
+    // false for required types regardless of input
+    return nullable_values() && valid_bits_.is_void();
   }
 
-  // Throw away levels from start_levels_position to levels_position_.
-  // Will update levels_position_, levels_written_, and levels_capacity_
-  // accordingly and move the levels to left to fill in the gap.
-  // It will resize the buffer without releasing the memory allocation.
-  void ThrowAwayLevels(int64_t start_levels_position) {
-    ARROW_DCHECK_LE(levels_position_, levels_written_);
-    ARROW_DCHECK_LE(start_levels_position, levels_position_);
-    ARROW_DCHECK_GT(this->max_def_level(), 0);
-    ARROW_DCHECK_NE(def_levels_, nullptr);
-
-    int64_t gap = levels_position_ - start_levels_position;
-    if (gap == 0) return;
-
-    int64_t levels_remaining = levels_written_ - gap;
+  uint8_t* values() const final { return 
reinterpret_cast<uint8_t*>(value_sink_.data()); }
 
-    auto left_shift = [&](::arrow::ResizableBuffer* buffer) {
-      auto* data = buffer->mutable_data_as<int16_t>();
-      std::copy(data + levels_position_, data + levels_written_,
-                data + start_levels_position);
-      PARQUET_THROW_NOT_OK(buffer->Resize(levels_remaining * sizeof(int16_t),
-                                          /*shrink_to_fit=*/false));
-    };
+  int64_t values_written() const final { return value_sink_.values_count(); }
 
-    left_shift(def_levels_.get());
-
-    if (this->max_rep_level() > 0) {
-      ARROW_DCHECK_NE(rep_levels_, nullptr);
-      left_shift(rep_levels_.get());
-    }
-
-    levels_written_ -= gap;
-    levels_position_ -= gap;
-    levels_capacity_ -= gap;
+  const void* ReadDictionary(int32_t* dictionary_length) final {
+    return reinterpret_cast<const 
void*>(Base::ReadDictionary(dictionary_length));
   }
 
+  int64_t ReadRecords(int64_t num_records) override;
+
   // Skip records that we have in our buffer. This function is only for
   // non-repeated fields.
-  int64_t SkipRecordsInBufferNonRepeated(int64_t num_records) {
-    ARROW_DCHECK_EQ(this->max_rep_level(), 0);
-    if (!this->has_values_to_process() || num_records == 0) return 0;
-
-    int64_t remaining_records = levels_written_ - levels_position_;
-    int64_t skipped_records = std::min(num_records, remaining_records);
-    int64_t start_levels_position = levels_position_;
-    // Since there is no repetition, number of levels equals number of records.
-    levels_position_ += skipped_records;
-
-    // We skipped the levels by incrementing 'levels_position_'. For values
-    // we do not have a buffer, so we need to read them and throw them away.
-    // First we need to figure out how many present/not-null values there are.
-    int64_t values_to_read =
-        std::count(def_levels() + start_levels_position, def_levels() + 
levels_position_,
-                   this->max_def_level());
-
-    // Now that we have figured out number of values to read, we do not need
-    // these levels anymore. We will remove these values from the buffer.
-    // This requires shifting the levels in the buffer to left. So this will
-    // update levels_position_ and levels_written_.
-    ThrowAwayLevels(start_levels_position);
-    // For values, we do not have them in buffer, so we will read them and
-    // throw them away.
-    ReadAndThrowAwayValues(values_to_read);
-
-    // Mark the levels as read in the underlying column reader.
-    this->ConsumeBufferedValues(skipped_records);
-
-    return skipped_records;
-  }
+  int64_t SkipRecordsInBufferNonRepeated(int64_t num_records);
 
   // Attempts to skip num_records from the buffer. Will throw away levels
   // and corresponding values for the records it skipped and consumes them 
from the
   // underlying decoder. Will advance levels_position_ and update
   // at_record_start_.
   // Returns how many records were skipped.
-  int64_t DelimitAndSkipRecordsInBuffer(int64_t num_records) {
-    if (num_records == 0) return 0;
-    // Look at the buffered levels, delimit them based on
-    // (rep_level == 0), report back how many records are in there, and
-    // fill in how many not-null values (def_level == max_def_level_).
-    // DelimitRecords updates levels_position_.
-    int64_t start_levels_position = levels_position_;
-    int64_t values_seen = 0;
-    int64_t skipped_records = DelimitRecords(num_records, &values_seen);
-    ReadAndThrowAwayValues(values_seen);
-    // Mark those levels and values as consumed in the underlying page.
-    // This must be done before we throw away levels since it updates
-    // levels_position_ and levels_written_.
-    this->ConsumeBufferedValues(levels_position_ - start_levels_position);
-    // Updated levels_position_ and levels_written_.
-    ThrowAwayLevels(start_levels_position);
-    return skipped_records;
-  }
+  int64_t DelimitAndSkipRecordsInBuffer(int64_t num_records);
 
   // Skip records for repeated fields. For repeated fields, we are technically
   // reading and throwing away the levels and values since we do not know the 
record
   // boundaries in advance. Keep filling the buffer and skipping until we 
reach the
   // desired number of records or we run out of values in the column chunk.
   // Returns number of skipped records.
-  int64_t SkipRecordsRepeated(int64_t num_records) {
-    ARROW_DCHECK_GT(this->max_rep_level(), 0);
-    int64_t skipped_records = 0;
-
-    // First consume what is in the buffer.
-    if (levels_position_ < levels_written_) {
-      // This updates at_record_start_.
-      skipped_records = DelimitAndSkipRecordsInBuffer(num_records);
-    }
+  int64_t SkipRecordsRepeated(int64_t num_records);
 
-    int64_t level_batch_size =
-        std::max<int64_t>(kMinLevelBatchSize, num_records - skipped_records);
-
-    // If 'at_record_start_' is false, but (skipped_records == num_records), it
-    // means that for the last record that was counted, we have not seen all
-    // of its values yet.
-    while (!at_record_start_ || skipped_records < num_records) {
-      // Is there more data to read in this row group?
-      // HasNextInternal() will advance to the next page if necessary.
-      if (!this->HasNextInternal()) {
-        if (!at_record_start_) {
-          // We ended the row group while inside a record that we haven't seen
-          // the end of yet. So increment the record count for the last record
-          // in the row group
-          ++skipped_records;
-          at_record_start_ = true;
-        }
-        break;
-      }
+  // Skip 'num_values' values from the current page.
+  void SkipValuesInPage(int64_t num_values);
 
-      // Read some more levels.
-      int64_t batch_size =
-          std::min(level_batch_size, this->available_values_current_page());
-      // No more data in column. This must be an empty page.
-      // If we had exhausted the last page, HasNextInternal() must have 
advanced
-      // to the next page. So there must be available values to process.
-      if (batch_size == 0) {
-        break;
+  int64_t SkipRecords(int64_t num_records) override;
+
+  // We may outwardly have the appearance of having exhausted a column chunk
+  // when in fact we are in the middle of processing the last batch
+  bool has_buffered_levels() const {
+    ARROW_DCHECK_LE(levels_position_, levels_written());
+    return levels_position_ < levels_written();
+  }
+
+  std::shared_ptr<ResizableBuffer> ReleaseValues() override {
+    return value_sink_.ReleaseValues(this->pool_);
+  }
+
+  std::shared_ptr<ResizableBuffer> ReleaseIsValid() final {
+    return valid_bits_.ReleaseValues(this->pool_);  // nullptr if void
+  }
+
+  // Process written repetition/definition levels to reach the end of
+  // records. Only used for repeated fields.
+  // Process no more levels than necessary to delimit the indicated
+  // number of logical records. Updates internal state of RecordReader
+  //
+  // \return Number of records delimited
+  int64_t DelimitRecords(int64_t num_records, int64_t* values_seen);
+
+  void Reserve(int64_t capacity) override;
+
+  void Reset() override;
+
+  void SetPageReader(std::unique_ptr<PageReader> reader) override;
+
+  bool HasMoreData() const override {
+    return Base::HasPageReader();  // Surprising legacy behaviour
+  }
+
+  const ColumnDescriptor* descr() const override { return this->descr_; }
+
+  // Reads repeated records from the buffered levels and returns number of 
records
+  // read. Fills in values_to_read and null_count.
+  int64_t ReadRepeatedRecordsInBuffer(int64_t num_records, int64_t* 
values_to_read,
+                                      int64_t* null_count);
+
+  // Reads optional records from the buffered levels and returns number of 
records
+  // read. Fills in values_to_read and null_count.
+  int64_t ReadOptionalRecordsInBuffer(int64_t num_records, int64_t* 
values_to_read,
+                                      int64_t* null_count);
+
+  // Reads dense for optional records. First it figures out how many values to
+  // read.
+  void ReadDenseForOptionalInBuffer(int64_t start_levels_position,
+                                    int64_t* values_to_read);
+
+  // Reads spaced for optional or repeated fields.
+  void ReadSpacedForOptionalOrRepeatedInBuffer(int64_t start_levels_position,
+                                               int64_t* values_to_read,
+                                               int64_t* null_count);
+
+  // Read records from the buffered levels.
+  // Return number of logical records read.
+  // Updates levels_position_, values_written_, and null_count_.
+  int64_t ReadRecordDataInBuffer(int64_t num_records);
+
+  void DebugPrintState() override;
+
+ protected:
+  auto value_sink() -> ValueSink& { return value_sink_; }
+  auto value_sink() const -> const ValueSink& { return value_sink_; }
+
+ private:
+  ValueSink value_sink_;
+  /// \brief Each bit corresponds to one element in 'values_' and specifies if 
it
+  /// is null or not null.
+  /// Not set if leaf type is not nullable or read_dense_for_nullable is true.
+  ValiditySinkBuffer valid_bits_ = ValiditySinkBuffer::MakeVoid();
+  LevelInfo leaf_info_;
+
+  /// \brief Buffer for definition levels.
+  ///
+  /// May contain eagerly decoded levels as required to figure out record 
boundaries for
+  /// repeated fields. `level_position_` is the number of level processed for 
the current
+  /// decoded values. For flat required fields, `def_levels_` and 
`rep_levels_` are
+  /// not populated nor allocated.
+  /// `def_levels_` and `rep_levels_` must be of the same size if present.
+  LevelSinkBuffer def_levels_ = LevelSinkBuffer::MakeVoid();
+  /// \brief Buffer for repetition levels. Only populated for repeated fields.
+  LevelSinkBuffer rep_levels_ = LevelSinkBuffer::MakeVoid();
+  /// \brief Position of the next level that should be consumed.
+  int64_t levels_position_ = 0;
+
+  int64_t null_count_ = 0;
+
+  bool at_record_start_ = true;
+};
+
+/**************************************
+ *  TypedRecordReader Implementation  *
+ **************************************/
+
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::ReadRecords(int64_t num_records) {
+  if (num_records == 0) return 0;
+  // Delimit records, then read values at the end
+  int64_t records_read = 0;
+
+  if (has_buffered_levels()) {
+    records_read += ReadRecordDataInBuffer(num_records);
+  }
+
+  int64_t level_batch_size = std::max<int64_t>(kMinLevelBatchSize, 
num_records);
+
+  // If we are in the middle of a record, we continue until reaching the
+  // desired number of records or the end of the current record if we've found
+  // enough records
+  while (!at_record_start_ || records_read < num_records) {
+    // Is there more data to read in this row group?
+    if (!this->EnsureDataPage()) {
+      if (!at_record_start_) {
+        // We ended the row group while inside a record that we haven't seen
+        // the end of yet. So increment the record count for the last record in
+        // the row group
+        ++records_read;
+        at_record_start_ = true;
       }
+      break;
+    }
 
-      // For skipping we will read the levels and append them to the end
-      // of the def_levels and rep_levels just like for read.
-      ReserveLevels(batch_size);
+    /// We perform multiple batch reads until we either exhaust the row group
+    /// or observe the desired number of records
+    const int32_t batch_size =
+        narrow_min(level_batch_size, this->available_values_current_page());
 
-      int16_t* def_levels = this->def_levels() + levels_written_;
-      int16_t* rep_levels = this->rep_levels() + levels_written_;
+    // No more data in column
+    if (batch_size == 0) {
+      break;
+    }
 
-      if (this->ReadDefinitionLevels(batch_size, def_levels) != batch_size) {
-        throw ParquetException(kErrorRepDefLevelNotMatchesNumValues);
+    if (this->max_def_level() > 0) {
+      def_levels_.Decode(this->def_levels_decoder_, batch_size);
+      if (this->max_rep_level() > 0) {
+        rep_levels_.Decode(this->rep_levels_decoder_, batch_size);
       }
-      if (this->ReadRepetitionLevels(batch_size, rep_levels) != batch_size) {
-        throw ParquetException(kErrorRepDefLevelNotMatchesNumValues);
+      records_read += ReadRecordDataInBuffer(num_records - records_read);
+    } else {
+      // No repetition and definition levels, we can read values directly
+      const auto count = narrow_min(num_records - records_read, batch_size);
+      records_read += ReadRecordDataInBuffer(count);
+    }
+  }
+
+  return records_read;
+}
+
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::SkipRecordsInBufferNonRepeated(
+    int64_t num_records) {
+  ARROW_DCHECK_EQ(this->max_rep_level(), 0);
+  if (!this->has_buffered_levels() || num_records == 0) return 0;
+
+  const int64_t remaining_records_64 = levels_written() - levels_position_;
+  ARROW_DCHECK_LE(remaining_records_64, std::numeric_limits<int32_t>::max());
+  const auto remaining_records = static_cast<int32_t>(remaining_records_64);
+  const int32_t skipped_records = narrow_min(num_records, remaining_records);
+  const int64_t remaining_levels_pos = levels_position_ + skipped_records;
+
+  // We skipped the levels by incrementing 'levels_position_'. For values
+  // we do not have a buffer, so we need to read them and throw them away.
+  // First we need to figure out how many present/not-null values there are.
+  const auto values_to_read = static_cast<int32_t>(
+      std::count(def_levels() + levels_position_, def_levels() + 
remaining_levels_pos,
+                 this->max_def_level()));
+
+  // Now that we have figured out number of values to read, we do not need
+  // these levels anymore. We will remove these values from the buffer.
+  def_levels_.Erase(levels_position_, remaining_levels_pos);
+
+  // For values, we do not have them in buffer, so we will read them and
+  // throw them away.
+  SkipValuesInPage(values_to_read);
+
+  // Mark the levels as read in the underlying column reader.
+  this->MarkValuesAsConsumed(skipped_records);
+
+  return skipped_records;
+}
+
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::DelimitAndSkipRecordsInBuffer(
+    int64_t num_records) {
+  if (num_records == 0) return 0;
+  // Look at the buffered levels, delimit them based on
+  // (rep_level == 0), report back how many records are in there, and
+  // fill in how many not-null values (def_level == max_def_level_).
+  // DelimitRecords updates levels_position_.
+  int64_t start_levels_position = levels_position_;
+  int64_t values_seen = 0;
+  int64_t skipped_records = DelimitRecords(num_records, &values_seen);
+  SkipValuesInPage(values_seen);
+  // Mark those levels and values as consumed in the underlying page.
+  // This must be done before we throw away levels since it updates
+  // levels_position_ and levels_written().
+  this->MarkValuesAsConsumed(clamp_to<int32_t>(levels_position_ - 
start_levels_position));
+  // Updated levels_position_ and levels_written().
+  def_levels_.Erase(start_levels_position, levels_position_);
+  rep_levels_.Erase(start_levels_position, levels_position_);
+  levels_position_ = start_levels_position;
+  return skipped_records;
+}
+
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::SkipRecordsRepeated(int64_t 
num_records) {
+  ARROW_DCHECK_GT(this->max_rep_level(), 0);
+  int64_t skipped_records = 0;
+
+  // First consume what is in the buffer.
+  if (has_buffered_levels()) {
+    // This updates at_record_start_.
+    skipped_records = DelimitAndSkipRecordsInBuffer(num_records);
+  }
+
+  int64_t level_batch_size =
+      std::max<int64_t>(kMinLevelBatchSize, num_records - skipped_records);
+
+  // If 'at_record_start_' is false, but (skipped_records == num_records), it
+  // means that for the last record that was counted, we have not seen all
+  // of its values yet.
+  while (!at_record_start_ || skipped_records < num_records) {
+    // Is there more data to read in this row group?
+    // HasNextInternal() will advance to the next page if necessary.
+    if (!this->EnsureDataPage()) {
+      if (!at_record_start_) {
+        // We ended the row group while inside a record that we haven't seen
+        // the end of yet. So increment the record count for the last record
+        // in the row group
+        ++skipped_records;
+        at_record_start_ = true;
       }
+      break;
+    }
+
+    // Read some more levels.
+    const int64_t batch_size_64 =
+        std::min<int64_t>(level_batch_size, 
this->available_values_current_page());
+    // available_values_current_page fits in int32_t
+    const auto batch_size = static_cast<int32_t>(batch_size_64);
 
-      levels_written_ += batch_size;
-      int64_t remaining_records = num_records - skipped_records;
-      // This updates at_record_start_.
-      skipped_records += DelimitAndSkipRecordsInBuffer(remaining_records);
+    // No more data in column. This must be an empty page.
+    // If we had exhausted the last page, HasNextInternal() must have advanced
+    // to the next page. So there must be available values to process.
+    if (batch_size == 0) {
+      break;
     }
 
+    def_levels_.Decode(this->def_levels_decoder_, batch_size);
+    rep_levels_.Decode(this->rep_levels_decoder_, batch_size);
+    const int64_t remaining_records = num_records - skipped_records;
+    // This updates at_record_start_.
+    skipped_records += DelimitAndSkipRecordsInBuffer(remaining_records);
+  }
+
+  return skipped_records;
+}
+
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, kDic>::SkipValuesInPage(int64_t num_values) {
+  const int64_t values_read = this->current_decoder_.Skip(num_values);
+  if (values_read < num_values) {
+    std::stringstream ss;
+    ss << "Could not read and throw away " << num_values << " values";
+    throw ParquetException(ss.str());
+  }
+}
+
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::SkipRecords(int64_t num_records) {
+  if (num_records == 0) return 0;
+
+  // Top level required field. Number of records equals number of levels,
+  // and there is no read-ahead for levels.
+  if (this->max_rep_level() == 0 && this->max_def_level() == 0) {
+    return this->Skip(num_records);
+  } else if (this->max_rep_level() == 0) {
+    // Non-repeated optional field.
+    // First consume whatever is in the buffer.
+    int64_t skipped_records = SkipRecordsInBufferNonRepeated(num_records);
+    ARROW_DCHECK_LE(skipped_records, num_records);
+
+    // For records that we have not buffered, we will use the column
+    // reader's Skip to do the remaining Skip. Since the field is not
+    // repeated number of levels to skip is the same as number of records
+    // to skip.
+    skipped_records += this->Skip(num_records - skipped_records);
     return skipped_records;
   }
+  return this->SkipRecordsRepeated(num_records);
+}
 
-  // Read 'num_values' values and throw them away.
-  // Throws an error if it could not read 'num_values'.
-  void ReadAndThrowAwayValues(int64_t num_values) {
-    const int64_t values_read = this->current_decoder_.Skip(num_values);
-    if (values_read < num_values) {
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::DelimitRecords(int64_t num_records,
+                                                        int64_t* values_seen) {
+  if (ARROW_PREDICT_FALSE(num_records == 0 || !has_buffered_levels())) {
+    *values_seen = 0;
+    return 0;
+  }
+  int64_t records_read = 0;
+  const int16_t* const rep_levels = this->rep_levels();
+  const int16_t* const def_levels = this->def_levels();
+  ARROW_DCHECK_GT(this->max_rep_level(), 0);
+  // If at_record_start_ is true, we are seeing the start of a record
+  // for the second time, such as after repeated calls to
+  // DelimitRecords. In this case we must continue until we find
+  // another record start or exhaust the ColumnChunk
+  int64_t level = levels_position_;
+  if (at_record_start_) {
+    if (ARROW_PREDICT_FALSE(rep_levels[levels_position_] != 0)) {
       std::stringstream ss;
-      ss << "Could not read and throw away " << num_values << " values";
+      ss << "The repetition level at the start of a record must be 0 but got "
+         << rep_levels[levels_position_];
       throw ParquetException(ss.str());
     }
-  }
+    ++levels_position_;
+    // We have decided to consume the level at this position; therefore we
+    // must advance until we find another record boundary
+    at_record_start_ = false;
+  }
+
+  // Count logical records and number of non-null values to read
+  ARROW_DCHECK(!at_record_start_);
+  // Scan repetition levels to find record end
+  while (has_buffered_levels()) {
+    // We use an estimated batch size to simplify branching and
+    // improve performance in the common case. This might slow
+    // things down a bit if a single long record remains, though.
+    const int64_t stride =
+        std::min(levels_written() - levels_position_, num_records - 
records_read);
+    const int64_t position_end = levels_position_ + stride;
+    for (int64_t i = levels_position_; i < position_end; ++i) {
+      records_read += rep_levels[i] == 0;
+    }
+    levels_position_ = position_end;
+    if (records_read == num_records) {
+      // Check last rep_level reaches the boundary and
+      // pop the last level.
+      ARROW_CHECK_EQ(rep_levels[levels_position_ - 1], 0);
+      --levels_position_;
+      // We've found the number of records we were looking for. Set
+      // at_record_start_ to true and break
+      at_record_start_ = true;
+      break;
+    }
+  }
+  // Scan definition levels to find number of physical values
+  *values_seen = std::count(def_levels + level, def_levels + levels_position_,
+                            this->max_def_level());
+  return records_read;
+}
 
-  int64_t SkipRecords(int64_t num_records) override {
-    if (num_records == 0) return 0;
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, kDic>::Reserve(int64_t extra_values) {
+  value_sink_.ReserveValues(extra_values);
+  valid_bits_.ReserveValues(extra_values);  // potentially no-op if void
+  def_levels_.ReserveValues(extra_values);  // potentially no-op if void
+  rep_levels_.ReserveValues(extra_values);  // potentially no-op if void
+}
 
-    // Top level required field. Number of records equals to number of levels,
-    // and there is not read-ahead for levels.
-    if (this->max_rep_level() == 0 && this->max_def_level() == 0) {
-      return this->Skip(num_records);
-    }
-    int64_t skipped_records = 0;
-    if (this->max_rep_level() == 0) {
-      // Non-repeated optional field.
-      // First consume whatever is in the buffer.
-      skipped_records = SkipRecordsInBufferNonRepeated(num_records);
-
-      ARROW_DCHECK_LE(skipped_records, num_records);
-
-      // For records that we have not buffered, we will use the column
-      // reader's Skip to do the remaining Skip. Since the field is not
-      // repeated number of levels to skip is the same as number of records
-      // to skip.
-      skipped_records += this->Skip(num_records - skipped_records);
-    } else {
-      skipped_records += this->SkipRecordsRepeated(num_records);
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, kDic>::Reset() {
+  null_count_ = 0;
+  value_sink_.ResetValues();
+  valid_bits_.ResetValues();  // potentially no-op if void
+  // Must keep eagerly decoded levels
+  def_levels_.Erase(0, levels_position_);  // potentially no-op if void
+  rep_levels_.Erase(0, levels_position_);  // potentially no-op if void
+  levels_position_ = 0;
+}
+
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, 
kDic>::SetPageReader(std::unique_ptr<PageReader> reader) {
+  at_record_start_ = true;
+  Base::SetPageReader(std::move(reader));
+  // At most one dictionary in Parquet column chunk and it has to be the first 
page.
+  if (this->HasPageReader() && this->EnsureDataPage()) {
+    if (auto* dict = this->current_dict_decoder()) {
+      value_sink_.OnNewDictionary(*dict);
     }
-    return skipped_records;
   }
+}
 
-  // We may outwardly have the appearance of having exhausted a column chunk
-  // when in fact we are in the middle of processing the last batch
-  bool has_values_to_process() const { return levels_position_ < 
levels_written_; }
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::ReadRepeatedRecordsInBuffer(
+    int64_t num_records, int64_t* values_to_read, int64_t* null_count) {
+  const int64_t start_levels_position = levels_position_;
+  // Note that repeated records may be required or nullable. If they have
+  // an optional parent in the path, they will be nullable, otherwise,
+  // they are required. We use leaf_info_->HasNullableValues() that looks
+  // at repeated_ancestor_def_level to determine if it is required or
+  // nullable. Even if they are required, we may have to read ahead and
+  // delimit the records to get the right number of values and they will
+  // have associated levels.
+  int64_t records_read = DelimitRecords(num_records, values_to_read);
+  if (valid_bits_.is_void()) {  // not nullable or read_dense_for_nullable
+    // This is only reading in the current page so this fits in an int32.
+    value_sink_.ReadValuesDense(*this->current_decoder_.get(),
+                                clamp_to<int32_t>(*values_to_read));
+    // null_count is always 0 for required.
+    ARROW_DCHECK_EQ(*null_count, 0);
+  } else {
+    ReadSpacedForOptionalOrRepeatedInBuffer(start_levels_position, 
values_to_read,
+                                            null_count);
+  }
+  return records_read;
+}
 
-  std::shared_ptr<ResizableBuffer> ReleaseValues() override {
-    if (uses_values_) {
-      auto result = values_;
-      PARQUET_THROW_NOT_OK(
-          result->Resize(bytes_for_values(values_written_), 
/*shrink_to_fit=*/true));
-      values_ = AllocateBuffer(this->pool_);
-      values_capacity_ = 0;
-      return result;
-    } else {
-      return nullptr;
-    }
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::ReadOptionalRecordsInBuffer(
+    int64_t num_records, int64_t* values_to_read, int64_t* null_count) {
+  const int64_t start_levels_position = levels_position_;
+  // No repetition levels, skip delimiting logic. Each level represents a
+  // null or not null entry
+  const int64_t records_read =
+      std::min<int64_t>(levels_written() - levels_position_, num_records);
+  // This is advanced by DelimitRecords for the repeated field case above.
+  levels_position_ += records_read;
+
+  // Optional fields are always nullable.
+  if (read_dense_for_nullable()) {
+    ReadDenseForOptionalInBuffer(start_levels_position, values_to_read);
+    // We don't need to update null_count when reading dense. It should be
+    // already set to 0.
+    ARROW_DCHECK_EQ(*null_count, 0);
+  } else {
+    ReadSpacedForOptionalOrRepeatedInBuffer(start_levels_position, 
values_to_read,
+                                            null_count);
   }
+  return records_read;
+}
 
-  std::shared_ptr<ResizableBuffer> ReleaseIsValid() override {
-    if (nullable_values()) {
-      auto result = valid_bits_;
-      
PARQUET_THROW_NOT_OK(result->Resize(bit_util::BytesForBits(values_written_),
-                                          /*shrink_to_fit=*/true));
-      valid_bits_ = AllocateBuffer(this->pool_);
-      return result;
-    } else {
-      return nullptr;
-    }
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, kDic>::ReadDenseForOptionalInBuffer(
+    int64_t start_levels_position, int64_t* values_to_read) {
+  // levels_position_ must already be incremented based on number of records
+  // read.
+  ARROW_DCHECK_GE(levels_position_, start_levels_position);
+
+  // When reading dense we need to figure out number of values to read.
+  const int16_t* def_levels = this->def_levels();
+  *values_to_read += std::count(def_levels + start_levels_position,
+                                def_levels + levels_position_, 
this->max_def_level());
+  // This is only reading in the current page so this fits in an int32.
+  value_sink_.ReadValuesDense(*this->current_decoder_.get(),
+                              clamp_to<int32_t>(*values_to_read));
+}
+
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, kDic>::ReadSpacedForOptionalOrRepeatedInBuffer(
+    int64_t start_levels_position, int64_t* values_to_read, int64_t* 
null_count) {
+  // levels_position_ must already be incremented based on number of records
+  // read.
+  const int64_t valid_bits_offset = valid_bits_.values_count();
+  const auto result =
+      valid_bits_.ReadFromDefLevels(def_levels() + start_levels_position,
+                                    levels_position_ - start_levels_position, 
leaf_info_);
+
+  *values_to_read = result.values_read - result.null_count;
+  *null_count = result.null_count;
+
+  // This is only reading in the current page so this fits in an int32.
+  value_sink_.ReadValuesSpaced(*this->current_decoder_.get(),
+                               clamp_to<int32_t>(result.values_read),
+                               clamp_to<int32_t>(*null_count),
+                               /* valid_bits= */ valid_bits_.data(),
+                               /* valid_bits_offset= */ valid_bits_offset);
+}
+
+template <typename DT, typename VS, bool kDic>
+int64_t TypedRecordReader<DT, VS, kDic>::ReadRecordDataInBuffer(int64_t 
num_records) {
+  // The value and validity sinks reserve their own capacity as they read, so
+  // there is no need to pre-reserve an upper bound here.
+  const int64_t start_levels_position = levels_position_;
+
+  // To be updated by the function calls below for each of the repetition
+  // types.
+  int64_t records_read = 0;
+  int64_t values_to_read = 0;
+  int64_t null_count = 0;
+  if (this->max_rep_level() > 0) {
+    // Repeated fields may be nullable or not.
+    // This call updates levels_position_.
+    records_read = ReadRepeatedRecordsInBuffer(num_records, &values_to_read, 
&null_count);
+  } else if (this->max_def_level() > 0) {
+    // Non-repeated optional values are always nullable.
+    // This call updates levels_position_.
+    ARROW_DCHECK(nullable_values());
+    records_read = ReadOptionalRecordsInBuffer(num_records, &values_to_read, 
&null_count);
+  } else {
+    ARROW_DCHECK(!nullable_values());
+    values_to_read = num_records;
+    // This is only reading in the current page so this fits in an int32.
+    value_sink_.ReadValuesDense(*this->current_decoder_.get(),
+                                clamp_to<int32_t>(values_to_read));
+    records_read = num_records;
+    // We don't need to update null_count, since it is 0.
+  }
+
+  ARROW_DCHECK_GE(records_read, 0);
+  ARROW_DCHECK_GE(values_to_read, 0);
+  ARROW_DCHECK_GE(null_count, 0);
+
+  // The values have already been accounted for by the value sink.
+  if (read_dense_for_nullable()) {
+    ARROW_DCHECK_EQ(null_count, 0);
+  } else {
+    null_count_ += null_count;
+  }
+  // Total values, including null spaces, if any
+  if (this->max_def_level() > 0) {
+    // Optional, repeated, or some mix thereof
+    // This is only reading in the current page so this fits in an int32.
+    this->MarkValuesAsConsumed(
+        clamp_to<int32_t>(levels_position_ - start_levels_position));
+  } else {
+    // Flat, non-repeated
+    // This is only reading in the current page so this fits in an int32.
+    this->MarkValuesAsConsumed(clamp_to<int32_t>(values_to_read));
   }
 
-  // Process written repetition/definition levels to reach the end of
-  // records. Only used for repeated fields.
-  // Process no more levels than necessary to delimit the indicated
-  // number of logical records. Updates internal state of RecordReader
-  //
-  // \return Number of records delimited
-  int64_t DelimitRecords(int64_t num_records, int64_t* values_seen) {
-    if (ARROW_PREDICT_FALSE(num_records == 0 || levels_position_ == 
levels_written_)) {
-      *values_seen = 0;
-      return 0;
-    }
-    int64_t records_read = 0;
-    const int16_t* const rep_levels = this->rep_levels();
-    const int16_t* const def_levels = this->def_levels();
-    ARROW_DCHECK_GT(this->max_rep_level(), 0);
-    // If at_record_start_ is true, we are seeing the start of a record
-    // for the second time, such as after repeated calls to
-    // DelimitRecords. In this case we must continue until we find
-    // another record start or exhausting the ColumnChunk
-    int64_t level = levels_position_;
-    if (at_record_start_) {
-      if (ARROW_PREDICT_FALSE(rep_levels[levels_position_] != 0)) {
-        std::stringstream ss;
-        ss << "The repetition level at the start of a record must be 0 but got 
"
-           << rep_levels[levels_position_];
-        throw ParquetException(ss.str());
-      }
-      ++levels_position_;
-      // We have decided to consume the level at this position; therefore we
-      // must advance until we find another record boundary
-      at_record_start_ = false;
+  return records_read;
+}
+
+template <typename DT, typename VS, bool kDic>
+void TypedRecordReader<DT, VS, kDic>::DebugPrintState() {
+  const int16_t* def_levels = this->def_levels();
+  const int16_t* rep_levels = this->rep_levels();
+  const int64_t total_levels_read = levels_position_;
+
+  if (leaf_info_.def_level > 0) {
+    std::cout << "def levels: ";
+    for (int64_t i = 0; i < total_levels_read; ++i) {
+      std::cout << def_levels[i] << " ";
     }
+    std::cout << std::endl;
+  }
 
-    // Count logical records and number of non-null values to read
-    ARROW_DCHECK(!at_record_start_);
-    // Scan repetition levels to find record end
-    while (levels_position_ < levels_written_) {
-      // We use an estimated batch size to simplify branching and
-      // improve performance in the common case. This might slow
-      // things down a bit if a single long record remains, though.
-      int64_t stride =
-          std::min(levels_written_ - levels_position_, num_records - 
records_read);
-      const int64_t position_end = levels_position_ + stride;
-      for (int64_t i = levels_position_; i < position_end; ++i) {
-        records_read += rep_levels[i] == 0;
-      }
-      levels_position_ = position_end;
-      if (records_read == num_records) {
-        // Check last rep_level reaches the boundary and
-        // pop the last level.
-        ARROW_CHECK_EQ(rep_levels[levels_position_ - 1], 0);
-        --levels_position_;
-        // We've found the number of records we were looking for. Set
-        // at_record_start_ to true and break
-        at_record_start_ = true;
-        break;
-      }
+  if (leaf_info_.rep_level > 0) {
+    std::cout << "rep levels: ";
+    for (int64_t i = 0; i < total_levels_read; ++i) {
+      std::cout << rep_levels[i] << " ";
     }
-    // Scan definition levels to find number of physical values
-    *values_seen = std::count(def_levels + level, def_levels + 
levels_position_,
-                              this->max_def_level());
-    return records_read;
+    std::cout << std::endl;
   }
 
-  void Reserve(int64_t capacity) override {
-    ReserveLevels(capacity);
-    ReserveValues(capacity);
+  std::cout << "values: ";
+  value_sink_.DebugPrintState();
+  std::cout << std::endl;
+}
+
+/*******************************
+ *  RequiredTypedRecordReader  *
+ *******************************/
+
+template <typename D>
+struct RequiredTypedRecordReaderTraits {
+  using DType = D;
+  using DefLevelDecoder = NewLevelDecoder;
+  using RepLevelDecoder = NewLevelDecoder;
+};
+
+/// A special type of record reader for required data.
+///
+/// Definition and repetition levels are all null in this case and the data 
encoded
+/// correspond directly to the
+template <typename DType, typename ValueSink = ValueSinkBuffer<typename 
DType::c_type>,
+          bool kReadDictionary = false>
+class RequiredTypedRecordReader
+    : public ColumnChunkReader<RequiredTypedRecordReaderTraits<DType>>,
+      virtual public RecordReader {
+ public:
+  using T = typename DType::c_type;
+  using Base = ColumnChunkReader<RequiredTypedRecordReaderTraits<DType>>;
+
+  RequiredTypedRecordReader(const ColumnDescriptor* descr, MemoryPool* pool,
+                            ValueSink value_sink)
+      : Base(descr, pool, NewLevelDecoder(descr->max_definition_level()),
+             NewLevelDecoder(descr->max_repetition_level())),
+        value_sink_(std::move(value_sink)) {
+    RequiredTypedRecordReader::Reset();
+    ARROW_DCHECK_EQ(descr->max_definition_level(), 0);
+    ARROW_DCHECK_EQ(descr->max_repetition_level(), 0);
   }
 
-  int64_t UpdateCapacity(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);
+  uint8_t* values() const final { return 
reinterpret_cast<uint8_t*>(value_sink_.data()); }
+
+  int64_t values_written() const final { return value_sink_.values_count(); }
+
+  int16_t* def_levels() const final { return nullptr; }
+
+  int16_t* rep_levels() const final { return nullptr; }
+
+  int64_t levels_position() const final { return 0; }
+
+  int64_t levels_written() const final { return 0; }
+
+  int64_t null_count() const final { return 0; }
+
+  bool nullable_values() const final { return false; }
+
+  bool read_dictionary() const final { return kReadDictionary; }
+
+  bool read_dense_for_nullable() const final { return false; }
+
+  const void* ReadDictionary(int32_t* dictionary_length) final {
+    return reinterpret_cast<const 
void*>(Base::ReadDictionary(dictionary_length));
   }
 
-  void ReserveLevels(int64_t extra_levels) {
-    if (this->max_def_level() > 0) {
-      const int64_t new_levels_capacity =
-          UpdateCapacity(levels_capacity_, levels_written_, extra_levels);
-      if (new_levels_capacity > levels_capacity_) {
-        constexpr auto kItemSize = static_cast<int64_t>(sizeof(int16_t));
-        int64_t capacity_in_bytes = -1;
-        if (MultiplyWithOverflow(new_levels_capacity, kItemSize, 
&capacity_in_bytes)) {
-          throw ParquetException("Allocation size too large (corrupt file?)");
-        }
-        PARQUET_THROW_NOT_OK(
-            def_levels_->Resize(capacity_in_bytes, /*shrink_to_fit=*/false));
-        if (this->max_rep_level() > 0) {
-          PARQUET_THROW_NOT_OK(
-              rep_levels_->Resize(capacity_in_bytes, /*shrink_to_fit=*/false));
-        }
-        levels_capacity_ = new_levels_capacity;
-      }
-    }
+  int64_t ReadRecords(int64_t num_records) final;
+
+  int64_t SkipRecords(int64_t num_records) final { return 
this->Skip(num_records); }
+
+  std::shared_ptr<ResizableBuffer> ReleaseValues() final {
+    return value_sink_.ReleaseValues(this->pool_);
   }
 
-  virtual void ReserveValues(int64_t extra_values) {
-    const int64_t new_values_capacity =
-        UpdateCapacity(values_capacity_, values_written_, extra_values);
-    if (new_values_capacity > values_capacity_) {
-      // XXX(wesm): A hack to avoid memory allocation when reading directly
-      // into builder classes
-      if (uses_values_) {
-        
PARQUET_THROW_NOT_OK(values_->Resize(bytes_for_values(new_values_capacity),
-                                             /*shrink_to_fit=*/false));
-      }
-      values_capacity_ = new_values_capacity;
-    }
-    if (nullable_values() && !read_dense_for_nullable_) {
-      int64_t valid_bytes_new = bit_util::BytesForBits(values_capacity_);
-      if (valid_bits_->size() < valid_bytes_new) {
-        int64_t valid_bytes_old = bit_util::BytesForBits(values_written_);
-        PARQUET_THROW_NOT_OK(
-            valid_bits_->Resize(valid_bytes_new, /*shrink_to_fit=*/false));
-
-        // Avoid valgrind warnings
-        memset(valid_bits_->mutable_data() + valid_bytes_old, 0,
-               static_cast<size_t>(valid_bytes_new - valid_bytes_old));
+  std::shared_ptr<ResizableBuffer> ReleaseIsValid() final { return nullptr; }
+
+  void Reserve(int64_t extra_values) final { ReserveValues(extra_values); }
+
+  void Reset() final;
+
+  void SetPageReader(std::unique_ptr<PageReader> reader) final {
+    Base::SetPageReader(std::move(reader));
+    // At most one dictionary in Parquet column chunk and it has to be the 
first page.
+    if (this->HasPageReader() && this->EnsureDataPage()) {
+      if (auto* dict = this->current_dict_decoder()) {
+        value_sink_.OnNewDictionary(*dict);
       }
     }
   }
 
-  void Reset() override {
-    ResetValues();
+  bool HasMoreData() const override {
+    return Base::HasPageReader();  // Surprising legacy behaviour
+  }
 
-    if (levels_written_ > 0) {
-      // Throw away levels from 0 to levels_position_.
-      ThrowAwayLevels(0);
-    }
+  const ColumnDescriptor* descr() const final { return this->descr_; }
+
+  void DebugPrintState() final;
 
-    // Call Finish on the binary builders to reset them
+ protected:
+  auto value_sink() -> ValueSink& { return value_sink_; }
+  auto value_sink() const -> const ValueSink& { return value_sink_; }
+
+ private:
+  ValueSink value_sink_;
+
+  void ReserveValues(int64_t extra_values) { 
value_sink_.ReserveValues(extra_values); }
+};
+
+/**********************************************
+ *  RequiredTypedRecordReader Implementation  *
+ **********************************************/
+
+template <typename DT, typename VS, bool kDic>
+int64_t RequiredTypedRecordReader<DT, VS, kDic>::ReadRecords(int64_t 
num_records) {
+  if (num_records <= 0) {
+    return 0;
   }
 
-  void SetPageReader(std::unique_ptr<PageReader> reader) override {
-    at_record_start_ = true;
-    this->pager_ = std::move(reader);
-    ResetDecoders();
+  int64_t records_read = 0;
+  do {
+    // Is there more data to read in this row group?
+    if (!this->EnsureDataPage()) {
+      break;
+    }
+
+    const int32_t batch_size =
+        narrow_min(num_records - records_read, 
this->available_values_current_page());
+    value_sink_.ReadValuesDense(*this->current_decoder_.get(), batch_size);
+    this->MarkValuesAsConsumed(batch_size);
+
+    records_read += batch_size;
+  } while (records_read < num_records);
+
+  return records_read;
+}
+
+template <typename DT, typename VS, bool kDic>
+void RequiredTypedRecordReader<DT, VS, kDic>::Reset() {
+  if (values_written() > 0) {
+    value_sink_.ResetValues();
   }
+}
+
+template <typename DT, typename VS, bool kDic>
+void RequiredTypedRecordReader<DT, VS, kDic>::DebugPrintState() {
+  std::cout << "values: ";
+  value_sink_.DebugPrintState();
+  std::cout << std::endl;
+}
 
-  bool HasMoreData() const override { return this->pager_ != nullptr; }
+/***********************************
+ *  FlatOptionalTypedRecordReader  *
+ ***********************************/
 
-  const ColumnDescriptor* descr() const override { return this->descr_; }
+template <typename DT>
+struct FlatOptionalTypedRecordReaderTraits {
+  using DType = DT;
+  using DefLevelDecoder = PageLevelToBitmapDecoder;
+  using RepLevelDecoder = NewLevelDecoder;
+};
 
-  // Dictionary decoders must be reset when advancing row groups
-  void ResetDecoders() { this->decoders_.clear(); }
-
-  virtual void ReadValuesSpaced(int64_t values_with_nulls, int64_t null_count) 
{
-    uint8_t* valid_bits = valid_bits_->mutable_data();
-    const int64_t valid_bits_offset = values_written_;
-
-    int64_t num_decoded = this->current_decoder_->DecodeSpaced(
-        ValuesHead<T>(), static_cast<int>(values_with_nulls),
-        static_cast<int>(null_count), valid_bits, valid_bits_offset);
-    CheckNumberDecoded(num_decoded, values_with_nulls);
-  }
-
-  virtual void ReadValuesDense(int64_t values_to_read) {
-    int64_t num_decoded =
-        this->current_decoder_->Decode(ValuesHead<T>(), 
static_cast<int>(values_to_read));
-    CheckNumberDecoded(num_decoded, values_to_read);
-  }
-
-  // Reads repeated records and returns number of records read. Fills in
-  // values_to_read and null_count.
-  int64_t ReadRepeatedRecords(int64_t num_records, int64_t* values_to_read,
-                              int64_t* null_count) {
-    const int64_t start_levels_position = levels_position_;
-    // Note that repeated records may be required or nullable. If they have
-    // an optional parent in the path, they will be nullable, otherwise,
-    // they are required. We use leaf_info_->HasNullableValues() that looks
-    // at repeated_ancestor_def_level to determine if it is required or
-    // nullable. Even if they are required, we may have to read ahead and
-    // delimit the records to get the right number of values and they will
-    // have associated levels.
-    int64_t records_read = DelimitRecords(num_records, values_to_read);
-    if (!nullable_values() || read_dense_for_nullable_) {
-      ReadValuesDense(*values_to_read);
-      // null_count is always 0 for required.
-      ARROW_DCHECK_EQ(*null_count, 0);
-    } else {
-      ReadSpacedForOptionalOrRepeated(start_levels_position, values_to_read, 
null_count);
-    }
-    return records_read;
-  }
-
-  // Reads optional records and returns number of records read. Fills in
-  // values_to_read and null_count.
-  int64_t ReadOptionalRecords(int64_t num_records, int64_t* values_to_read,
-                              int64_t* null_count) {
-    const int64_t start_levels_position = levels_position_;
-    // No repetition levels, skip delimiting logic. Each level represents a
-    // null or not null entry
-    int64_t records_read =
-        std::min<int64_t>(levels_written_ - levels_position_, num_records);
-    // This is advanced by DelimitRecords for the repeated field case above.
-    levels_position_ += records_read;
-
-    // Optional fields are always nullable.
-    if (read_dense_for_nullable_) {
-      ReadDenseForOptional(start_levels_position, values_to_read);
-      // We don't need to update null_count when reading dense. It should be
-      // already set to 0.
-      ARROW_DCHECK_EQ(*null_count, 0);
-    } else {
-      ReadSpacedForOptionalOrRepeated(start_levels_position, values_to_read, 
null_count);
+/// A specialized record reader for flat optional data.
+///
+/// In this special case, the max definition level is 1 and these correspond 
to the arrow
+/// array we are building. A special level decoder is used to bypass decoding 
completely
+/// and only copy the bitmap into the Arrow buffer.
+template <typename DType>
+class FlatOptionalTypedRecordReader
+    : public ColumnChunkReader<FlatOptionalTypedRecordReaderTraits<DType>>,
+      virtual public RecordReader {
+ public:
+  using T = typename DType::c_type;
+  using Base = ColumnChunkReader<FlatOptionalTypedRecordReaderTraits<DType>>;
+  using ValueSink = ValueSinkBuffer<T>;
+
+  FlatOptionalTypedRecordReader(const ColumnDescriptor* descr, MemoryPool* 
pool,
+                                bool read_dense_for_nullable, ValueSink 
value_sink)
+      : Base(descr, pool, PageLevelToBitmapDecoder(), NewLevelDecoder(0)),
+        value_sink_(std::move(value_sink)) {
+    ARROW_DCHECK_EQ(descr->max_definition_level(), 1);
+    ARROW_DCHECK_EQ(descr->max_repetition_level(), 0);
+    ARROW_DCHECK(descr->schema_node()->is_optional());
+    if (!read_dense_for_nullable) {
+      valid_bits_ = ValiditySinkBuffer::MakeAllocated(this->pool_);
     }
-    return records_read;
   }
 
-  // Reads required records and returns number of records read. Fills in
-  // values_to_read.
-  int64_t ReadRequiredRecords(int64_t num_records, int64_t* values_to_read) {
-    *values_to_read = num_records;
-    ReadValuesDense(*values_to_read);
-    return num_records;
+  uint8_t* values() const final { return 
reinterpret_cast<uint8_t*>(value_sink_.data()); }
+
+  int64_t values_written() const final { return value_sink_.values_count(); }
+
+  int16_t* def_levels() const final { return nullptr; }
+
+  int16_t* rep_levels() const final { return nullptr; }
+
+  int64_t levels_position() const final { return 0; }
+
+  int64_t levels_written() const final { return 0; }
+
+  int64_t null_count() const final { return null_count_; }
+
+  bool nullable_values() const final { return true; }
+
+  bool read_dictionary() const final { return false; }
+
+  bool read_dense_for_nullable() const final { return valid_bits_.is_void(); }
+
+  const void* ReadDictionary(int32_t* dictionary_length) final {
+    return reinterpret_cast<const 
void*>(Base::ReadDictionary(dictionary_length));
   }
 
-  // Reads dense for optional records. First it figures out how many values to
-  // read.
-  void ReadDenseForOptional(int64_t start_levels_position, int64_t* 
values_to_read) {
-    // levels_position_ must already be incremented based on number of records
-    // read.
-    ARROW_DCHECK_GE(levels_position_, start_levels_position);
+  int64_t ReadRecords(int64_t num_records) final;
 
-    // When reading dense we need to figure out number of values to read.
-    const int16_t* def_levels = this->def_levels();
-    *values_to_read += std::count(def_levels + start_levels_position,
-                                  def_levels + levels_position_, 
this->max_def_level());
-    ReadValuesDense(*values_to_read);
+  int64_t SkipRecords(int64_t num_records) final { return 
this->Skip(num_records); }
+
+  std::shared_ptr<ResizableBuffer> ReleaseValues() final {
+    return value_sink_.ReleaseValues(this->pool_);
   }
 
-  // Reads spaced for optional or repeated fields.
-  void ReadSpacedForOptionalOrRepeated(int64_t start_levels_position,
-                                       int64_t* values_to_read, int64_t* 
null_count) {
-    // levels_position_ must already be incremented based on number of records
-    // read.
-    ARROW_DCHECK_GE(levels_position_, start_levels_position);
-    ValidityBitmapInputOutput validity_io;
-    validity_io.values_read_upper_bound = levels_position_ - 
start_levels_position;
-    validity_io.valid_bits = valid_bits_->mutable_data();
-    validity_io.valid_bits_offset = values_written_;
-
-    DefLevelsToBitmap(def_levels() + start_levels_position,
-                      levels_position_ - start_levels_position, leaf_info_, 
&validity_io);
-    *values_to_read = validity_io.values_read - validity_io.null_count;
-    *null_count = validity_io.null_count;
-    ARROW_DCHECK_GE(*values_to_read, 0);
-    ARROW_DCHECK_GE(*null_count, 0);
-    ReadValuesSpaced(validity_io.values_read, *null_count);
+  std::shared_ptr<ResizableBuffer> ReleaseIsValid() final {
+    return valid_bits_.ReleaseValues(this->pool_);  // nullptr if void
   }
 
-  // Return number of logical records read.
-  // Updates levels_position_, values_written_, and null_count_.
-  int64_t ReadRecordData(int64_t num_records) {
-    // Conservative upper bound
-    const int64_t possible_num_values =
-        std::max<int64_t>(num_records, levels_written_ - levels_position_);
-    ReserveValues(static_cast<size_t>(possible_num_values));
-
-    const int64_t start_levels_position = levels_position_;
-
-    // To be updated by the function calls below for each of the repetition
-    // types.
-    int64_t records_read = 0;
-    int64_t values_to_read = 0;
-    int64_t null_count = 0;
-    if (this->max_rep_level() > 0) {
-      // Repeated fields may be nullable or not.
-      // This call updates levels_position_.
-      records_read = ReadRepeatedRecords(num_records, &values_to_read, 
&null_count);
-    } else if (this->max_def_level() > 0) {
-      // Non-repeated optional values are always nullable.
-      // This call updates levels_position_.
-      ARROW_DCHECK(nullable_values());
-      records_read = ReadOptionalRecords(num_records, &values_to_read, 
&null_count);
-    } else {
-      ARROW_DCHECK(!nullable_values());
-      records_read = ReadRequiredRecords(num_records, &values_to_read);
-      // We don't need to update null_count, since it is 0.
-    }
+  void Reserve(int64_t extra_values) final {
+    value_sink_.ReserveValues(extra_values);
+    valid_bits_.ReserveValues(extra_values);
+  }
 
-    ARROW_DCHECK_GE(records_read, 0);
-    ARROW_DCHECK_GE(values_to_read, 0);
-    ARROW_DCHECK_GE(null_count, 0);
+  void Reset() final;
 
-    if (read_dense_for_nullable_) {
-      values_written_ += values_to_read;
-      ARROW_DCHECK_EQ(null_count, 0);
-    } else {
-      values_written_ += values_to_read + null_count;
-      null_count_ += null_count;
-    }
-    // Total values, including null spaces, if any
-    if (this->max_def_level() > 0) {
-      // Optional, repeated, or some mix thereof
-      this->ConsumeBufferedValues(levels_position_ - start_levels_position);
-    } else {
-      // Flat, non-repeated
-      this->ConsumeBufferedValues(values_to_read);
+  void SetPageReader(std::unique_ptr<PageReader> reader) final {
+    Base::SetPageReader(std::move(reader));
+    // At most one dictionary in Parquet column chunk and it has to be the 
first page.
+    if (this->HasPageReader() && this->EnsureDataPage()) {
+      if (auto* dict = this->current_dict_decoder()) {
+        value_sink_.OnNewDictionary(*dict);
+      }
     }
+  }
 
-    return records_read;
+  bool HasMoreData() const override {
+    return Base::HasPageReader();  // Surprising legacy behaviour
   }
 
-  void DebugPrintState() override {
-    const int16_t* def_levels = this->def_levels();
-    const int16_t* rep_levels = this->rep_levels();
-    const int64_t total_levels_read = levels_position_;
+  const ColumnDescriptor* descr() const final { return this->descr_; }
 
-    const T* vals = reinterpret_cast<const T*>(this->values());
+  void DebugPrintState() final;
 
-    if (leaf_info_.def_level > 0) {
-      std::cout << "def levels: ";
-      for (int64_t i = 0; i < total_levels_read; ++i) {
-        std::cout << def_levels[i] << " ";
-      }
-      std::cout << std::endl;
-    }
+ protected:
+  auto value_sink() -> ValueSink& { return value_sink_; }
+  auto value_sink() const -> const ValueSink& { return value_sink_; }
 
-    if (leaf_info_.rep_level > 0) {
-      std::cout << "rep levels: ";
-      for (int64_t i = 0; i < total_levels_read; ++i) {
-        std::cout << rep_levels[i] << " ";
-      }
-      std::cout << std::endl;
+ private:
+  ValueSink value_sink_;
+  ValiditySinkBuffer valid_bits_ = ValiditySinkBuffer::MakeVoid();
+  int64_t null_count_ = 0;
+};
+
+/**************************************************
+ *  FlatOptionalTypedRecordReader Implementation  *
+ **************************************************/
+
+template <typename DT>
+int64_t FlatOptionalTypedRecordReader<DT>::ReadRecords(int64_t num_records) {
+  if (num_records <= 0) {
+    return 0;
+  }
+
+  int64_t records_read = 0;
+
+  do {
+    // Is there more data to read in this row group?
+    if (!this->EnsureDataPage()) {
+      break;
     }
 
-    std::cout << "values: ";
-    for (int64_t i = 0; i < this->values_written(); ++i) {
-      std::cout << vals[i] << " ";
+    const int32_t batch_size =
+        narrow_min(num_records - records_read, 
this->available_values_current_page());
+    if (read_dense_for_nullable()) {
+      const auto result = this->def_levels_decoder_.CountUpTo(true, 
batch_size);
+      value_sink_.ReadValuesDense(*this->current_decoder_.get(), 
result.matching_count);
+      records_read += result.processed_count;

Review Comment:
   The dense fast path does not verify that `CountUpTo` processed the requested 
batch. A truncated or malformed level stream can return fewer levels; this 
advances `records_read` by less than `batch_size` while the page still reports 
the same values available, and a subsequent zero-progress call can loop forever 
instead of reporting the corrupt page. Check `result.processed_count == 
batch_size` before decoding values, as the spaced path does.



-- 
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