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

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


The following commit(s) were added to refs/heads/main by this push:
     new 43d04c2f6f GH-51496: [C++] Fix 32-bit block offset overflows in 
SwissTable (#51497)
43d04c2f6f is described below

commit 43d04c2f6fad8c9c820ffc608ce3606f534b0bde
Author: Jáchym Barvínek <[email protected]>
AuthorDate: Mon Sep 28 06:05:10 2026 +0200

    GH-51496: [C++] Fix 32-bit block offset overflows in SwissTable (#51497)
    
    ### Rationale for this change
    
    `SwissTable::early_filter_imp_avx2_x8` computes each block's byte offset 
with `_mm256_mullo_epi32`, so modulo 2^32. With 32-bit group ids a block takes 
40 bytes, and every block with id >= 107,374,183 (ceil(2^32 / 40)) gets a 
wrapped offset. The kernel then reads the status bytes from the wrong place, 
reports no match (or the wrong slot), and `find()` misses the key. There is no 
error.
    
    A hash join whose build side has more than ~403M distinct keys 
(`log_blocks` >= 27) loses about 20%, 60% and 80% of its matches at 
`log_blocks` 27, 28 and 29. The kernel is only dispatched on CPUs with AVX2 and 
efficient BMI2 (`HasEfficientBmi2()`, that is Intel), so the same join is 
correct on AMD, on ARM and with `ARROW_USER_SIMD_LEVEL=NONE`. See GH-XXXXX for 
the full analysis and a pyarrow reproducer.
    
    While checking the other block offset computations I found the same 
overflow in scalar code. When `grow_double()` reinserts an overflow entry and 
the entry's new home block is full, it probes the next blocks at 
`blocks_new->mutable_data() + block_id_new * block_size_after`, a 32-bit 
multiply. GH-45506 fixed the line just above but not this one. Once a single 
table grows to 2^27 blocks, such an entry is written to an unrelated location 
and can't be found. This happens on any CPU. A mul [...]
    
    ### What changes are included in this PR?
    
    - `key_map_internal_avx2.cc`: compute `voffset_A`/`voffset_B` in 64 bits 
with `_mm256_mul_epu32` on the even and odd 32-bit block ids. This is the same 
even/odd split as before, now without wraparound, and matches what 
`extract_group_ids_avx2` does since GH-44513.
    - `key_map_internal.cc`: use `mutable_block_data()`, which computes the 
offset in 64 bits, for the probe in `grow_double()`.
    - New `key_map_test.cc` (added to `arrow-compute-row-test` in CMake and 
Meson) with two large-memory tests:
      - `SwissTable.EarlyFilterOver4GB`: a table with 2^27 blocks and no hash 
array (5.4 GB). It inserts keys into the blocks on both sides of the 4GB 
boundary and checks `early_filter` for every supported hardware flag set 
(scalar, and AVX2 where available).
      - `SwissTable.GrowOver4GB` (about 14.5 GB peak): it grows a table from 
2^26 to 2^27 blocks with an entry that has to move past a full block beyond 
4GB, then checks `early_filter` and `find` for all keys. It triggers the growth 
through `num_inserted()` and the 75% fill threshold, and asserts that the table 
grew.
    
    ### Are these changes tested?
    
    Yes. Both new tests are `LARGE_MEMORY_TEST`s, so they only run with 
`ARROW_LARGE_MEMORY_TESTS=ON`. The AVX2 part of `EarlyFilterOver4GB` only 
exercises the bug on Intel CPUs with AVX2 and BMI2.
    
    The results below are from Google Cloud `n1-highmem-32` (Intel Xeon @ 
2.00GHz, AVX2, BMI2 and AVX-512), Ubuntu 24.04, gcc 13.3, with 
`ARROW_RUNTIME_SIMD_LEVEL=MAX`, in Release and in Debug with 
`BUILD_WARNING_LEVEL=CHECKIN`:
    
    - `main` with only the new tests added:
      - `EarlyFilterOver4GB` fails at `key_map_test.cc:78` with 
`local_slots[i]` 7, expected 0, at `block_id = 107374183`, `hardware_flags = 
32` (AVX2). It passes with `ARROW_USER_SIMD_LEVEL=NONE`.
      - `GrowOver4GB` fails at `key_map_test.cc:151` (key 8 not found, 
`hardware_flags = 0`), with any SIMD level.
    - The early filter fix alone: `EarlyFilterOver4GB` passes, and 
`GrowOver4GB` still fails at `:151`.
    - This PR: both tests pass, also with `ARROW_USER_SIMD_LEVEL=NONE` (about 
10 s and 16 s in Release).
      - `arrow-compute-row-test`: 90 tests pass (the 88 existing ones and the 2 
new ones).
      - `arrow-acero-hash-join-node-test`: 36 pass and 1 is skipped 
(`BuildSideLargeRowIds`, which is skipped in its body), the same as on `main`.
      - No new compiler warnings in the Debug `-Werror` build.
    
    Before the fix, the pyarrow reproducer from the issue gave these results on 
the same machine type with pyarrow 25.0.1 (an inner join where all 2,000,000 
probe keys have a match):
    
    | build rows | rows returned | with `ARROW_USER_SIMD_LEVEL=NONE` |
    |---:|---:|---:|
    | 390,000,000 | 2,000,000 (100.00%) | 2,000,000 |
    | 450,000,000 | 1,603,163 (80.16%) | 2,000,000 |
    | 900,000,000 | 814,637 (40.73%) | 2,000,000 |
    
    I did not rerun the pyarrow reproducer on a patched build. The C++ tests 
above cover the same code path.
    
    ### Are there any user-facing changes?
    
    No API changes. Hash joins with more than ~403M distinct build keys on 
Intel AVX2 CPUs, and single hash tables that grow past ~403M keys, now return 
correct results.
    
    **This PR contains a "Critical Fix".** Both bugs silently produce incorrect 
data. The AVX2 early filter makes hash joins with more than ~403M distinct 
build-side keys drop 20% to 80% of the matching rows on Intel CPUs with AVX2 
and BMI2. The `grow_double()` overflow can misplace entries, and so lose keys, 
when a hash table grows to 2^27 blocks or more (past ~403M keys), on any CPU.
    
    ### Was AI used for this PR?
    
    In accordance to the [AI generation 
guidelines](https://arrow.apache.org/docs/dev/developers/overview.html#ai-generated-code),
 please disclose below whether and how AI was used in this PR.
    
    **PR code and description written by:**
    
    - [ ] Human
    - [x] AI
    
    **Reviewed before submission by:**
    
    - [x] Human
    - [x] AI
    - [ ] Not reviewed
    
    Claude Code (Claude Opus) did the root-cause analysis and wrote the fix, 
the tests and this description, under my direction. I found the bug in a 
production pipeline. A separate AI pass reviewed the diff and rebuilt and reran 
the tests on an Intel machine.
    * GitHub Issue: #51496
    
    Authored-by: Jachym.Barvinek <[email protected]>
    Signed-off-by: Rossi Sun <[email protected]>
---
 cpp/src/arrow/compute/CMakeLists.txt           |   1 +
 cpp/src/arrow/compute/key_map_internal.cc      |   3 +-
 cpp/src/arrow/compute/key_map_internal_avx2.cc |   8 +-
 cpp/src/arrow/compute/key_map_test.cc          | 158 +++++++++++++++++++++++++
 cpp/src/arrow/compute/meson.build              |   1 +
 5 files changed, 166 insertions(+), 5 deletions(-)

diff --git a/cpp/src/arrow/compute/CMakeLists.txt 
b/cpp/src/arrow/compute/CMakeLists.txt
index 88f9ae9006..bfddb66358 100644
--- a/cpp/src/arrow/compute/CMakeLists.txt
+++ b/cpp/src/arrow/compute/CMakeLists.txt
@@ -175,6 +175,7 @@ add_arrow_compute_test(expression_test
 add_arrow_compute_test(row_test
                        SOURCES
                        key_hash_test.cc
+                       key_map_test.cc
                        light_array_test.cc
                        row/compare_test.cc
                        row/grouper_test.cc
diff --git a/cpp/src/arrow/compute/key_map_internal.cc 
b/cpp/src/arrow/compute/key_map_internal.cc
index 353449cf16..975a42c7d3 100644
--- a/cpp/src/arrow/compute/key_map_internal.cc
+++ b/cpp/src/arrow/compute/key_map_internal.cc
@@ -744,7 +744,8 @@ Status SwissTable::grow_double() {
           static_cast<int>(std::countl_zero(block_new & kHighBitOfEachByte) >> 
3);
       while (full_slots_new == kSlotsPerBlock) {
         block_id_new = (block_id_new + 1) & ((1 << log_blocks_after) - 1);
-        block_base_new = blocks_new->mutable_data() + block_id_new * 
block_size_after;
+        block_base_new = mutable_block_data(blocks_new->mutable_data(), 
block_id_new,
+                                            block_size_after);
         block_new = util::SafeLoadAs<uint64_t>(block_base_new);
         full_slots_new =
             static_cast<int>(std::countl_zero(block_new & kHighBitOfEachByte) 
>> 3);
diff --git a/cpp/src/arrow/compute/key_map_internal_avx2.cc 
b/cpp/src/arrow/compute/key_map_internal_avx2.cc
index 353d5a59e6..6f5e34c27b 100644
--- a/cpp/src/arrow/compute/key_map_internal_avx2.cc
+++ b/cpp/src/arrow/compute/key_map_internal_avx2.cc
@@ -52,12 +52,12 @@ int SwissTable::early_filter_imp_avx2_x8(const int 
num_hashes, const uint32_t* h
 
     // We now split inputs and process 4 at a time,
     // in order to process 64-bit blocks
+    // Block offsets are computed in 64 bits, as they may not fit in 32 bits.
     //
-    __m256i vblock_offset =
-        _mm256_mullo_epi32(vblock_id, _mm256_set1_epi32(num_block_bytes));
-    __m256i voffset_A = _mm256_and_si256(vblock_offset, 
_mm256_set1_epi64x(0xffffffff));
+    __m256i voffset_A = _mm256_mul_epu32(vblock_id, 
_mm256_set1_epi32(num_block_bytes));
     __m256i vstamp_A = _mm256_and_si256(vstamp, 
_mm256_set1_epi64x(0xffffffff));
-    __m256i voffset_B = _mm256_srli_epi64(vblock_offset, 32);
+    __m256i voffset_B = _mm256_mul_epu32(_mm256_srli_epi64(vblock_id, 32),
+                                         _mm256_set1_epi32(num_block_bytes));
     __m256i vstamp_B = _mm256_srli_epi64(vstamp, 32);
 
     auto blocks_i64 =
diff --git a/cpp/src/arrow/compute/key_map_test.cc 
b/cpp/src/arrow/compute/key_map_test.cc
new file mode 100644
index 0000000000..402d43f1b0
--- /dev/null
+++ b/cpp/src/arrow/compute/key_map_test.cc
@@ -0,0 +1,158 @@
+// 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.
+
+#include <gtest/gtest.h>
+
+#include <cstdint>
+#include <vector>
+
+#include "arrow/compute/key_map_internal.h"
+#include "arrow/memory_pool.h"
+#include "arrow/testing/gtest_util.h"
+#include "arrow/testing/util.h"
+#include "arrow/util/bit_util.h"
+#include "arrow/util/cpu_info.h"
+
+namespace arrow {
+
+using internal::CpuInfo;
+
+namespace compute {
+
+// With 32-bit group ids a block takes 40 bytes, so the byte offsets of the 
blocks from
+// id ceil(2^32 / 40) on do not fit in 32 bits.
+TEST(SwissTable, LARGE_MEMORY_TEST(EarlyFilterOver4GB)) {
+  if constexpr (sizeof(void*) == 4) {
+    GTEST_SKIP() << "Test only works on 64-bit platforms";
+  }
+
+  // 2^27 blocks of 40 bytes take 5GB.
+  constexpr int kLogBlocks = 27;
+  constexpr uint32_t kFirstBlockOver4GB = 107374183;
+  constexpr int kNumHashes = 16;
+
+  // One hash for each of the 8 blocks below and the 8 blocks above the 4GB 
boundary,
+  // all with a non-zero stamp.
+  std::vector<uint32_t> hashes(kNumHashes);
+  for (int i = 0; i < kNumHashes; ++i) {
+    uint32_t block_id = kFirstBlockOver4GB - kNumHashes / 2 + i;
+    hashes[i] = (block_id << (SwissTable::bits_hash_ - kLogBlocks)) | 1;
+  }
+
+  for (int64_t hardware_flags : GetSupportedHardwareFlags({CpuInfo::AVX2})) {
+    ARROW_SCOPED_TRACE("hardware_flags = ", hardware_flags);
+    SwissTable table;
+    ASSERT_OK(table.init(hardware_flags, default_memory_pool(), kLogBlocks,
+                         /*no_hash_array=*/true));
+    // Insert every other hash into the first slot of its block, leaving the 
other blocks
+    // empty.
+    for (int i = 0; i < kNumHashes; i += 2) {
+      uint32_t block_id = SwissTable::block_id_from_hash(hashes[i], 
kLogBlocks);
+      table.insert_into_empty_slot(SwissTable::global_slot_id(block_id, 0), 
hashes[i],
+                                   /*group_id=*/i);
+    }
+
+    uint8_t match_bitvector[kNumHashes / 8];
+    uint8_t local_slots[kNumHashes];
+    table.early_filter(kNumHashes, hashes.data(), match_bitvector, 
local_slots);
+    for (int i = 0; i < kNumHashes; ++i) {
+      ARROW_SCOPED_TRACE("block_id = ",
+                         SwissTable::block_id_from_hash(hashes[i], 
kLogBlocks));
+      // Inserted hashes match in the first slot, the others hit an empty 
block whose
+      // first slot is empty.
+      ASSERT_EQ(bit_util::GetBit(match_bitvector, i), i % 2 == 0);
+      ASSERT_EQ(local_slots[i], 0);
+    }
+  }
+}
+
+// When growing to a table over 4GB, entries that have to move past a full 
block must
+// still land in the right block.
+TEST(SwissTable, LARGE_MEMORY_TEST(GrowOver4GB)) {
+  if constexpr (sizeof(void*) == 4) {
+    GTEST_SKIP() << "Test only works on 64-bit platforms";
+  }
+
+  // Grow from 2^26 to 2^27 blocks.
+  constexpr int kLogBlocks = 26;
+  constexpr uint32_t kBlockId = (1u << kLogBlocks) - 2;
+  constexpr int kNumHashes = SwissTable::kSlotsPerBlock + 1;
+
+  // All these hashes map to block kBlockId before growing and to block 2 * 
kBlockId
+  // after. The first 8 fill these blocks, so the last one is in block 
kBlockId + 1 before
+  // growing and has to move to block 2 * kBlockId + 1, which is over 4GB.
+  std::vector<uint32_t> hashes(kNumHashes);
+  for (int i = 0; i < kNumHashes; ++i) {
+    hashes[i] = (kBlockId << (SwissTable::bits_hash_ - kLogBlocks)) | i;
+  }
+
+  // Key i is equal to group id i.
+  SwissTable::EqualImpl equal_impl =
+      [](int num_keys, const uint16_t* selection, const uint32_t* group_ids,
+         uint32_t* out_num_keys_mismatch, uint16_t* out_selection_mismatch, 
void*) {
+        *out_num_keys_mismatch = 0;
+        for (int i = 0; i < num_keys; ++i) {
+          uint16_t id = selection ? selection[i] : static_cast<uint16_t>(i);
+          if (group_ids[id] != id) {
+            out_selection_mismatch[(*out_num_keys_mismatch)++] = id;
+          }
+        }
+      };
+  SwissTable::AppendImpl append_impl = [](int, const uint16_t*, void*) {
+    return Status::OK();
+  };
+
+  for (int64_t hardware_flags : GetSupportedHardwareFlags({CpuInfo::AVX2})) {
+    ARROW_SCOPED_TRACE("hardware_flags = ", hardware_flags);
+    SwissTable table;
+    ASSERT_OK(table.init(hardware_flags, default_memory_pool(), kLogBlocks));
+    for (int i = 0; i < kNumHashes; ++i) {
+      uint32_t slot_id = SwissTable::global_slot_id(kBlockId, 0) + i;
+      table.insert_into_empty_slot(slot_id, hashes[i], /*group_id=*/i);
+      table.hashes()[slot_id] = hashes[i];
+    }
+
+    // The table grows when 75% of its slots are used. Pretend that it is one 
key short
+    // of that and insert one more key.
+    table.num_inserted(
+        static_cast<uint32_t>((int64_t{1} << (kLogBlocks + 3)) * 3 / 4 - 1));
+    util::TempVectorStack temp_stack;
+    ASSERT_OK(temp_stack.Init(default_memory_pool(), 64 * 
table.minibatch_size()));
+    uint16_t new_key_id = 0;
+    uint32_t new_key_hash = 0;
+    uint32_t new_group_id;
+    ASSERT_OK(table.map_new_keys(/*num_ids=*/1, &new_key_id, &new_key_hash, 
&new_group_id,
+                                 &temp_stack, equal_impl, append_impl,
+                                 /*callback_ctx=*/nullptr));
+    ASSERT_EQ(table.log_blocks(), kLogBlocks + 1);
+
+    uint8_t match_bitvector[(kNumHashes + 7) / 8];
+    uint8_t local_slots[kNumHashes];
+    uint32_t group_ids[kNumHashes];
+    table.early_filter(kNumHashes, hashes.data(), match_bitvector, 
local_slots);
+    table.find(kNumHashes, hashes.data(), match_bitvector, local_slots, 
group_ids,
+               &temp_stack, equal_impl, /*callback_ctx=*/nullptr);
+    for (int i = 0; i < kNumHashes; ++i) {
+      ARROW_SCOPED_TRACE("key = ", i);
+      ASSERT_TRUE(bit_util::GetBit(match_bitvector, i));
+      ASSERT_EQ(group_ids[i], i);
+    }
+  }
+}
+
+}  // namespace compute
+}  // namespace arrow
diff --git a/cpp/src/arrow/compute/meson.build 
b/cpp/src/arrow/compute/meson.build
index fed699e1ca..9641658a25 100644
--- a/cpp/src/arrow/compute/meson.build
+++ b/cpp/src/arrow/compute/meson.build
@@ -102,6 +102,7 @@ compute_tests = {
     'arrow-compute-row-test': {
         'sources': [
             'key_hash_test.cc',
+            'key_map_test.cc',
             'light_array_test.cc',
             'row/compare_test.cc',
             'row/grouper_test.cc',

Reply via email to