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