Skip to content

perf(deduplicator): compact Ray BTS edge buffers - #1078

Open
vertebrateqing wants to merge 2 commits into
datajuicer:mainfrom
vertebrateqing:codex/numpy-edge-buffer
Open

vertebrateqing wants to merge 2 commits into
datajuicer:mainfrom
vertebrateqing:codex/numpy-edge-buffer

Conversation

@vertebrateqing

@vertebrateqing vertebrateqing commented Sep 24, 2026 •

Copy link
Copy Markdown

Summary

  • replace per-destination Python list[tuple[int, int]] edge buffers with one packed NumPy array of signed-int64 edge pairs plus destination offsets on the common path; fall back to exact object-dtype storage only for UIDs outside the signed-int64 range
  • group outgoing edges by destination actor once, then transfer and consume each contiguous slice once
  • preallocate communication buffers and union received edges in bounded chunks
  • keep the public configuration and deduplication semantics unchanged

Part of #1031.

Correctness

The compact representation preserves the existing routing rules:

  • every edge is sent to the actor owning its source UID
  • a cross-partition edge is also sent to the actor owning its parent UID
  • union still selects the smaller root, so ordering within a destination does not change connected components
  • UIDs use compact signed int64 storage when representable; values outside that range fall back to object arrays without hashing, truncation, or collisions

Tests cover partition routing, one-time buffer consumption, signed and order-aligned parent arrays, connected components, out-of-range Python integer UIDs in both BTS communication phases, and cross-actor merging above 2**63.

Benchmark

Targeted edge-redistribution benchmark with 64 partitions. Each input parent relation produces two routed edges. Each result is the median of three fresh processes.

Environment: macOS 15.3, Apple M4, Python 3.12.14, NumPy 2.2.6, Ray 2.58.0.

Input relations main time PR time Speedup main peak RSS PR peak RSS RSS reduction
100K 0.0424 s 0.0092 s 4.61x 184.1 MiB 181.5 MiB 1.4%
300K 0.0811 s 0.0291 s 2.79x 236.3 MiB 219.5 MiB 7.1%
1M 0.2224 s 0.0997 s 2.23x 437.0 MiB 345.5 MiB 20.9%

At 1M input relations, peak RSS is reduced by 91.5 MiB. The incremental RSS slope from 100K to 1M is reduced by approximately 35.2%.

Validation

  • pre-commit run --all-files
  • standalone storage regression tests
  • Ray cross-actor connected-component test
  • Python Ray BTS MinHash English and Chinese end-to-end tests, including the UID path

The macOS build intentionally skips the optional C++ MinHash extensions, so the C++ implementation is left to Linux CI. This PR does not modify the C++ path.

@vertebrateqing

Copy link
Copy Markdown
Author

Hi maintainers, the unittest-partial workflow has been waiting for approval in the Testing environment since Sep 24. Could an environment reviewer approve the run when convenient? This PR is part of #1031. Thanks!

@Qirui-jiao

Copy link
Copy Markdown
Collaborator

Thanks a lot for this contribution! I have one suggestion before merging:

The PR says deduplication semantics are unchanged, but the UID range is now narrower. The old code stored UIDs as Python ints of any size. The new buffers are fixed to int64, so a UID ≥ 2^63 (for example, a uint64 or hash-derived __dj__uid passed to RayBTSMinhashDeduplicatorWithUid) now fails. Duplicate groups containing such UIDs deduplicate correctly on main but fail here. The failure only shows up in edge_redistribution, after the full MinHash pass, and communication() raises a bare OverflowError instead of the ValueError. #1031 asks us to keep error behavior unchanged and not add input limits, so please fall back to an object dtype or the old list path when a UID is out of int64 range.

Thanks again!

@vertebrateqing

Copy link
Copy Markdown
Author

Thanks @Qirui-jiao — addressed in 55c626c3.

The packed signed-int64 path remains the default. _parent_to_arrays() and communication() now retry with object-dtype buffers only when NumPy reports an integer overflow, so arbitrary Python integer UIDs (including values >= 2**63 and < -2**63) retain their exact values. The fallback routing uses Python integer arithmetic; it does not hash or truncate UIDs.

I added regression coverage for out-of-range values in both edge_redistribution() and communication(), plus a real two-actor Ray merge starting at 2**63. Local validation passed: 6 standalone storage tests, the cross-actor Ray test, Python Ray English/Chinese/WithUid end-to-end tests, and pre-commit run --all-files. The branch is also rebased onto the latest main.

This branch is waiting to be deployed

1 waiting deployment
Testing — 55c626c3 Waiting Sep 30, 2026 by vertebrateqing via unittest-dist #1776
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants