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

zhouyuan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new 03440927ec [VL] Sort shuffle row ID truncates offsets above 128 MiB 
(#13121)
03440927ec is described below

commit 03440927ec455a660f26aaa21ecc20ecfd0c6b1d
Author: jianzhenwu <[email protected]>
AuthorDate: Wed Sep 30 03:08:42 2026 +0800

    [VL] Sort shuffle row ID truncates offsets above 128 MiB (#13121)
---
 cpp/velox/shuffle/VeloxSortShuffleWriter.cc   | 79 +++++++++++++++++++++++----
 cpp/velox/shuffle/VeloxSortShuffleWriter.h    |  6 +-
 cpp/velox/tests/VeloxSortShuffleWriterTest.cc | 34 ++++++++++++
 3 files changed, 105 insertions(+), 14 deletions(-)

diff --git a/cpp/velox/shuffle/VeloxSortShuffleWriter.cc 
b/cpp/velox/shuffle/VeloxSortShuffleWriter.cc
index 0a31dd3d7d..49e88c0bf4 100644
--- a/cpp/velox/shuffle/VeloxSortShuffleWriter.cc
+++ b/cpp/velox/shuffle/VeloxSortShuffleWriter.cc
@@ -28,7 +28,9 @@ namespace gluten {
 namespace {
 
 constexpr uint32_t kMaskLower27Bits = (1 << 27) - 1;
+constexpr uint32_t kCompactRowIdOffsetLimit = kMaskLower27Bits + 1;
 constexpr uint64_t kMaskLower40Bits = (1UL << 40) - 1;
+constexpr uint32_t kCompactRowIdPageLimit = 1U << 13;
 constexpr uint32_t kPartitionIdStartByteIndex = 5;
 constexpr uint32_t kPartitionIdEndByteIndex = 7;
 constexpr uint32_t kMaxPageNumber = (1 << 13) - 1; // 13-bit max = 8191
@@ -187,6 +189,10 @@ arrow::Status VeloxSortShuffleWriter::insert(const 
facebook::velox::RowVectorPtr
     if (rows == 0) {
       auto minSizeRequired =
           (fixedRowSize_.has_value() ? fixedRowSize_.value() : 
rowSize_[rowOffset]) + sizeof(RowSizeType);
+      if (pages_.size() >= kCompactRowIdPageLimit) {
+        ARROW_RETURN_IF(offset_ == 0, arrow::Status::Invalid("Compact row ID 
page number exceeds its 13-bit limit."));
+        RETURN_NOT_OK(evictAllPartitions());
+      }
       acquireNewBuffer(static_cast<uint64_t>(memLimit), minSizeRequired);
       rows = maxRowsToInsert(rowOffset, remainingRows);
       ARROW_RETURN_IF(
@@ -196,13 +202,13 @@ arrow::Status VeloxSortShuffleWriter::insert(const 
facebook::velox::RowVectorPtr
     RETURN_NOT_OK(maybeSpill(rows));
     // Allocate newArray can trigger spill.
     growArrayIfNecessary(rows);
-    insertRows(row, rowOffset, rows);
+    RETURN_NOT_OK(insertRows(row, rowOffset, rows));
     rowOffset += rows;
   }
   return arrow::Status::OK();
 }
 
-void VeloxSortShuffleWriter::insertRows(
+arrow::Status VeloxSortShuffleWriter::insertRows(
     facebook::velox::row::CompactRow& compact,
     facebook::velox::vector_size_t offset,
     facebook::velox::vector_size_t size) {
@@ -211,14 +217,24 @@ void VeloxSortShuffleWriter::insertRows(
   for (auto i = 0; i < size; ++i) {
     auto row = offset + i;
     auto pid = row2Partition_[row];
+    ARROW_RETURN_IF(
+        pageNumber_ >= kCompactRowIdPageLimit,
+        arrow::Status::Invalid("Compact row ID page number exceeds its 13-bit 
limit."));
+    ARROW_RETURN_IF(
+        pageCursor_ >= kCompactRowIdOffsetLimit,
+        arrow::Status::Invalid("Compact row ID offset exceeds its 27-bit 
limit."));
+    auto recordSize = static_cast<uint64_t>(rowSize_[row]) + 
sizeof(RowSizeType);
+    ARROW_RETURN_IF(
+        pageCursor_ > currenPageSize_ || recordSize > currenPageSize_ - 
pageCursor_,
+        arrow::Status::Invalid("Compact row exceeds its shuffle page."));
     arrayPtr_[offset_++] = toCompactRowId(pid, pageNumber_, pageCursor_);
     // size(RowSize) | bytes
     memcpy(currentPage_ + pageCursor_, &rowSize_[row], sizeof(RowSizeType));
     offsets[i] = pageCursor_ + sizeof(RowSizeType);
-    pageCursor_ += rowSize_[row] + sizeof(RowSizeType);
-    VELOX_DCHECK_LE(pageCursor_, currenPageSize_);
+    pageCursor_ += recordSize;
   }
   compact.serialize(offset, size, offsets.data(), currentPage_);
+  return arrow::Status::OK();
 }
 
 arrow::Status VeloxSortShuffleWriter::maybeSpill(uint32_t nextRows) {
@@ -231,6 +247,14 @@ arrow::Status VeloxSortShuffleWriter::maybeSpill(uint32_t 
nextRows) {
 arrow::Status VeloxSortShuffleWriter::evictAllPartitions() {
   VELOX_CHECK(offset_ > 0);
   EvictGuard evictGuard{evictState_};
+  std::vector<uint32_t> pageSizes;
+  pageSizes.reserve(pages_.size());
+  for (const auto& page : pages_) {
+    pageSizes.push_back(page->size());
+  }
+  ARROW_RETURN_IF(
+      pageSizes.size() != pageAddresses_.size(),
+      arrow::Status::Invalid("Shuffle page addresses and buffers are out of 
sync."));
 
   auto numRecords = offset_;
   // offset_ is used for checking spillable data.
@@ -251,12 +275,12 @@ arrow::Status 
VeloxSortShuffleWriter::evictAllPartitions() {
   while (++cur < end) {
     auto curPid = extractPartitionId(arrayPtr_[cur]);
     if (curPid != pid) {
-      RETURN_NOT_OK(evictPartition(pid, begin, cur));
+      RETURN_NOT_OK(evictPartition(pid, begin, cur, pageSizes));
       pid = curPid;
       begin = cur;
     }
   }
-  RETURN_NOT_OK(evictPartition(pid, begin, cur));
+  RETURN_NOT_OK(evictPartition(pid, begin, cur, pageSizes));
 
   if (!stopped_) {
     // Preserve the last page for use.
@@ -281,8 +305,13 @@ arrow::Status VeloxSortShuffleWriter::evictAllPartitions() 
{
   return arrow::Status::OK();
 }
 
-arrow::Status VeloxSortShuffleWriter::evictPartition(uint32_t partitionId, 
size_t begin, size_t end) {
+arrow::Status VeloxSortShuffleWriter::evictPartition(
+    uint32_t partitionId,
+    size_t begin,
+    size_t end,
+    const std::vector<uint32_t>& pageSizes) {
   VELOX_DCHECK(begin < end);
+
   // Count copy row time into sortTime_.
   Timer sortTime{};
   // Serialize [begin, end)
@@ -293,8 +322,18 @@ arrow::Status 
VeloxSortShuffleWriter::evictPartition(uint32_t partitionId, size_
   auto index = begin;
   while (index < end) {
     auto pageIndex = extractPageNumberAndOffset(arrayPtr_[index]);
+    ARROW_RETURN_IF(
+        pageIndex.first >= pageAddresses_.size(), 
arrow::Status::Invalid("Compact row ID references an unknown page."));
+    auto pageSize = pageSizes[pageIndex.first];
+    ARROW_RETURN_IF(
+        pageSize < sizeof(RowSizeType) || pageIndex.second > pageSize - 
sizeof(RowSizeType),
+        arrow::Status::Invalid("Compact row ID offset is outside its shuffle 
page."));
     addr = pageAddresses_[pageIndex.first] + pageIndex.second;
-    recordSize = *(reinterpret_cast<RowSizeType*>(addr)) + sizeof(RowSizeType);
+    auto rowSize = *(reinterpret_cast<RowSizeType*>(addr));
+    ARROW_RETURN_IF(
+        rowSize > pageSize - pageIndex.second - sizeof(RowSizeType),
+        arrow::Status::Invalid("Compact row length exceeds its shuffle 
page."));
+    recordSize = rowSize + sizeof(RowSizeType);
     if (offset + recordSize > diskWriteBufferSize_ && offset > 0) {
       sortTime.stop();
       RETURN_NOT_OK(evictPartitionInternal(partitionId, index - begin, 
sortedBufferPtr_, offset));
@@ -355,16 +394,32 @@ facebook::velox::vector_size_t 
VeloxSortShuffleWriter::maxRowsToInsert(
   if (pages_.empty()) {
     return 0;
   }
+  if (pageCursor_ >= kCompactRowIdOffsetLimit) {
+    return 0;
+  }
   auto remainingBytes = pages_.back()->size() - pageCursor_;
   if (fixedRowSize_.has_value()) {
-    return std::min(
+    auto rowsThatFit = std::min(
         static_cast<facebook::velox::vector_size_t>(remainingBytes / 
(fixedRowSize_.value() + sizeof(RowSizeType))),
         remainingRows);
+    auto rowWidth = fixedRowSize_.value() + sizeof(RowSizeType);
+    auto rowsWithEncodableOffsets =
+        static_cast<facebook::velox::vector_size_t>((kMaskLower27Bits - 
pageCursor_) / rowWidth + 1);
+    return std::min(rowsThatFit, rowsWithEncodableOffsets);
   }
-  auto beginIter = rowSizePrefixSum_.begin() + 1 + offset;
   auto bytesWritten = rowSizePrefixSum_[offset];
-  auto iter = std::upper_bound(beginIter, rowSizePrefixSum_.end(), 
remainingBytes + bytesWritten);
-  return (facebook::velox::vector_size_t)(iter - beginIter);
+  auto pageFitBegin = rowSizePrefixSum_.begin() + 1 + offset;
+  auto pageFitEnd = std::upper_bound(pageFitBegin, rowSizePrefixSum_.end(), 
remainingBytes + bytesWritten);
+  auto rowsThatFit = static_cast<facebook::velox::vector_size_t>(pageFitEnd - 
pageFitBegin);
+
+  // Row IDs have room for offsets [0, 2^27). A row may extend past that 
boundary, but no later
+  // row may start there; start the next row on a fresh page instead.
+  auto rowStartBegin = rowSizePrefixSum_.begin() + offset;
+  auto rowStartEnd = rowSizePrefixSum_.end() - 1;
+  auto maxEncodableRowStart = bytesWritten + (kMaskLower27Bits - pageCursor_);
+  auto encodableEnd = std::upper_bound(rowStartBegin, rowStartEnd, 
maxEncodableRowStart);
+  auto rowsWithEncodableOffsets = 
static_cast<facebook::velox::vector_size_t>(encodableEnd - rowStartBegin);
+  return std::min(rowsThatFit, rowsWithEncodableOffsets);
 }
 
 void VeloxSortShuffleWriter::acquireNewBuffer(uint64_t memLimit, uint64_t 
minSizeRequired) {
diff --git a/cpp/velox/shuffle/VeloxSortShuffleWriter.h 
b/cpp/velox/shuffle/VeloxSortShuffleWriter.h
index 59541d2496..7aaaecab44 100644
--- a/cpp/velox/shuffle/VeloxSortShuffleWriter.h
+++ b/cpp/velox/shuffle/VeloxSortShuffleWriter.h
@@ -50,6 +50,8 @@ class VeloxSortShuffleWriter final : public 
VeloxShuffleWriter {
   int64_t totalC2RTime() const override;
 
  private:
+  friend class VeloxSortShuffleWriterTest;
+
   VeloxSortShuffleWriter(
       uint32_t numPartitions,
       const std::shared_ptr<PartitionWriter>& partitionWriter,
@@ -64,7 +66,7 @@ class VeloxSortShuffleWriter final : public 
VeloxShuffleWriter {
 
   arrow::Status insert(const facebook::velox::RowVectorPtr& vector, int64_t 
memLimit);
 
-  void insertRows(
+  arrow::Status insertRows(
       facebook::velox::row::CompactRow& compact,
       facebook::velox::vector_size_t offset,
       facebook::velox::vector_size_t size);
@@ -73,7 +75,7 @@ class VeloxSortShuffleWriter final : public 
VeloxShuffleWriter {
 
   arrow::Status evictAllPartitions();
 
-  arrow::Status evictPartition(uint32_t partitionId, size_t begin, size_t end);
+  arrow::Status evictPartition(uint32_t partitionId, size_t begin, size_t end, 
const std::vector<uint32_t>& pageSizes);
 
   arrow::Status evictPartitionInternal(uint32_t partitionId, uint32_t numRows, 
uint8_t* buffer, int64_t rawLength);
 
diff --git a/cpp/velox/tests/VeloxSortShuffleWriterTest.cc 
b/cpp/velox/tests/VeloxSortShuffleWriterTest.cc
index 130f1bb010..fcc0ef0ef7 100644
--- a/cpp/velox/tests/VeloxSortShuffleWriterTest.cc
+++ b/cpp/velox/tests/VeloxSortShuffleWriterTest.cc
@@ -45,6 +45,14 @@ class FakeBufferRssClient : public RssClient {
 
 class VeloxSortShuffleWriterTest : public VeloxShuffleWriterTestBase, public 
testing::Test {
  protected:
+  static size_t numPages(const VeloxSortShuffleWriter& writer) {
+    return writer.pageAddresses_.size();
+  }
+
+  static uint64_t rowIdAt(const VeloxSortShuffleWriter& writer, size_t index) {
+    return writer.arrayPtr_[index];
+  }
+
   static void SetUpTestSuite() {
     setUpVeloxBackend();
   }
@@ -97,4 +105,30 @@ TEST_F(VeloxSortShuffleWriterTest, pushCompleteRows) {
   EXPECT_EQ(10, rssClient->getReceiveTimes());
 }
 
+TEST_F(VeloxSortShuffleWriterTest, 
rollsOverPageBeforeCompactRowOffsetOverflows) {
+  constexpr uint32_t kCompactRowOffsetLimit = 1U << 27;
+  std::string largeValue(kCompactRowOffsetLimit + 1024 * 1024, 'x');
+  auto values = std::vector<StringView>{StringView(largeValue), 
StringView("small")};
+  auto rowVector = makeRowVector({makeFlatVector<StringView>(values)});
+  auto writeOptions = std::make_shared<SortShuffleWriterOptions>();
+  auto rssClient = std::make_shared<FakeBufferRssClient>();
+  auto shuffleWriter =
+      std::dynamic_pointer_cast<VeloxSortShuffleWriter>(createShuffleWriter(1, 
writeOptions, rssClient));
+
+  auto status = 
shuffleWriter->write(std::make_shared<VeloxColumnarBatch>(rowVector), 
ShuffleWriter::kMinMemLimit);
+  ASSERT_TRUE(status.ok()) << status.ToString();
+
+  // The second row must start on a new page because its offset would exceed 
the 27-bit row ID field.
+  ASSERT_EQ(2, numPages(*shuffleWriter));
+  constexpr uint64_t kOffsetMask = (1ULL << 27) - 1;
+  constexpr uint64_t kPageMask = ((1ULL << 40) - 1) >> 27;
+  EXPECT_EQ(0, rowIdAt(*shuffleWriter, 0) & kOffsetMask);
+  EXPECT_EQ(0, (rowIdAt(*shuffleWriter, 0) >> 27) & kPageMask);
+  EXPECT_EQ(0, rowIdAt(*shuffleWriter, 1) & kOffsetMask);
+  EXPECT_EQ(1, (rowIdAt(*shuffleWriter, 1) >> 27) & kPageMask);
+
+  ASSERT_TRUE(shuffleWriter->stop().ok());
+  EXPECT_EQ(2, rssClient->getReceiveTimes());
+}
+
 } // namespace gluten


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to