RapidsMPF: Đột phá tốc độ tráo đổi dữ liệu ngoài bộ nhớ GPU đạt 1,8 TiB/s
RapidsMPF là thư viện xử lý tráo đổi dữ liệu (shuffle) ngoài bộ nhớ (out-of-core) dành cho GPU, giúp giải quyết bài toán tràn bộ nhớ (OOM) khi thực hiện join, groupby hay sắp xếp trên tập dữ liệu lớn. Trên hệ thống DGX B200, thư viện này đạt thông lượng toàn cục khoảng 1,8 TiB/s và có khả năng tràn dữ liệu (spill) một cách có kiểm soát thay vì gây lỗi OOM.

RapidsMPF: Đột phá tốc độ tráo đổi dữ liệu ngoài bộ nhớ GPU đạt 1,8 TiB/s
RapidsMPF là một thư viện tráo đổi dữ liệu (shuffle) ngoài bộ nhớ (out-of-core) có thể tái sử dụng, biến những cơn đau đầu vì tràn bộ nhớ (OOM) thành một kế hoạch tràn dữ liệu (spill) có thể dự toán được. Trên hệ thống DGX B200 với 8 GPU Blackwell, thư viện này đạt thông lượng toàn cục khoảng 1,8 TiB/s khi dữ liệu còn nằm trong VRAM, và vẫn hoạt động ổn định ngay cả khi bị ép xuống giới hạn bộ nhớ chỉ 12 GiB mỗi GPU.
Vì sao shuffle là bài toán hóc búa?
Shuffle là khâu then chốt trong phân tích dữ liệu có cấu trúc, dù phân tán hay không. Đây là thành phần cốt lõi của các phép toán quan trọng như join, groupby, merge, sort. Một shuffle phân tán đầy đủ có thể phải di chuyển toàn bộ dữ liệu từ mọi tiến trình đến mọi tiến trình khác — tức là một phép all-to-all — và đây là thao tác cực kỳ tốn kém.
Điều đáng chú ý là shuffle không hề khó về mặt tính toán: việc băm (hash) khóa để định tuyến dữ liệu khá rẻ. Nhưng nó lại đắt đỏ trong một quy trình vì ba lý do:
- Ngốn bộ nhớ: shuffle có thể cần giữ một bản sao đầy đủ của toàn bộ dữ liệu, hoặc trong chế độ streaming, áp lực bộ nhớ tích tụ và gây OOM.
- Truyền tải: dữ liệu phải được di chuyển vật lý từ tiến trình A sang tiến trình B, hoặc từ node A sang node B, nên tốc độ bị giới hạn bởi tầng truyền tải.
- Đồng bộ hóa: dữ liệu đầu ra không thể được tiêu thụ cho đến khi tất cả các bên sản xuất hoàn tất việc đóng góp dữ liệu. Trong một engine bulk-synchronous, rào cản này có thể làm đình trệ toàn bộ kế hoạch thực thi.
Chính vì shuffle vừa khó, vừa chậm, vừa ngốn bộ nhớ lại vừa quan trọng, nên đây là điểm khởi đầu lịch sử của RapidsMPF.
Tại sao phép join cần shuffle?
Hãy cùng xem một ví dụ nhanh về join bảng. Giả sử có hai bảng partsupp và lineitem và ta muốn join chúng:
partsupp.join(
lineitem,
left_on=["ps_partkey", "ps_suppkey"],
right_on=["l_partkey", "l_suppkey"],
)
Trong trường hợp join trong bộ nhớ (in-memory), inner join gồm hai pha:
- Pha build: bảng nhỏ hơn (
partsupp) được quét và một bảng băm (hash table) được xây dựng trên các khóa join (ps_partkey,ps_suppkey), ánh xạ mỗi khóa đã băm tới dòng gốc. Bảng băm này phải được nạp đầy đủ trước khi pha probe bắt đầu. - Pha probe: bảng lớn hơn (
lineitem) được quét và khóa của mỗi dòng (l_partkey,l_suppkey) được băm. Các kết quả khớp giữa bảng build và bảng probe sẽ phát ra một dòng đầu ra kết hợp cột từ cả hai bảng. Những dòng không khớp sẽ bị loại bỏ.
Tối thiểu, phép join trong bộ nhớ này phải giữ ba bảng: bảng build, bảng probe, bảng đầu ra, cùng với bảng băm được xây dựng trên phía build.
Trong môi trường phân tán, điều này phức tạp hơn nhiều. Một tiến trình chỉ có thể join các dòng nằm trong bộ nhớ đang cư trú của nó. Để thực thi một distributed hash-join, ta phải:
- Quét bảng build, băm khóa join của từng dòng để chọn phân vùng đích:
hash(keys) % n_out_partitions, đóng gói và gửi mỗi dòng tới rank sẽ sở hữu nó. - Quét bảng probe và định tuyến các dòng theo cùng cách, để khóa probe hạ cánh xuống rank đang giữ khóa build có cùng giá trị băm.
- Chờ cho đến khi mọi rank gửi xong. Chỉ khi đó một rank mới được đảm bảo giữ mọi dòng từ cả hai bảng cho các khóa nó sở hữu.
- Chạy phép join trong bộ nhớ trên từng lát dữ liệu cục bộ của mỗi rank.
Sơ đồ các dòng dữ liệu cùng màu được shuffle vào cùng một phân vùng đầu ra trên các rank khác nhau
Trong trường hợp xấu nhất, nếu mọi giai đoạn đều được vật chất hóa đầy đủ trước khi giai đoạn tiếp theo bắt đầu, một rank duy nhất sẽ phải giữ đồng thời: bảng build, bảng probe, bảng build đã đóng gói, bảng probe đã đóng gói, lát build đã shuffle, lát probe đã shuffle, bảng băm trên lát build, và bảng đầu ra. Đó là lý do vì sao shuffle ngốn bộ nhớ chứ không ngốn tính toán.
RapidsMPF là gì?
RapidsMPF đã mở rộng kể từ ý tưởng ban đầu. Hiện nay nó là một thư viện gồm hai phần lớn:
- Thư viện shuffle được thiết kế cho việc tràn bộ nhớ / xử lý ngoài bộ nhớ, với tầng truyền tải được tăng tốc.
- Mạng actor (actor network) để xây dựng các pipeline dữ liệu streaming.
Người dùng hiện nay vẫn có thể chỉ áp dụng riêng thành phần shuffle của RapidsMPF (qua giao diện C++ hoặc Python). Chúng ta đã thấy sự áp dụng này trong NeMo-Curator và ở mức thử nghiệm trong Ray Data. Quan trọng nhất, cuDF Polars sử dụng RapidsMPF cho cả shuffle lẫn mạng actor.
Thư viện shuffle của RapidsMPF cần phải: nhanh, mở rộng được, hoạt động với dữ liệu lớn hơn VRAM (out-of-core), và có thể tái sử dụng.
Thiết lập đo điểm chuẩn
cuDF/RapidsMPF cung cấp một benchmark C++ dễ dùng tên là bench_shuffle, giúp nghiên cứu cách triển khai shuffle hoạt động trên nhiều cấu hình phần cứng và cấu hình khác nhau. Benchmark này tạo ra một lượng lớn số nguyên 32-bit ngẫu nhiên có thể tinh chỉnh cho mỗi rank (mỗi GPU), shuffle toàn bộ dữ liệu, rồi hoàn tất (không có join, chỉ shuffle).
Một số tham số chính được sử dụng trong bài đo:
-C ucxx: communicator sử dụng UCXX/UCX để bật truyền tải tăng tốc/GPUDirect RDMA.-c 10: 10 cột.-n 536870912: 536.870.912 dòng mỗi rank (tương đương 2 GiB mỗi cột với 4 byte/dòng).-p 1và-o 8: 1 phân vùng đầu vào mỗi rank, 8 phân vùng đầu ra (mỗi rank một phân vùng).-m pool: dùng RMM memory pool.-l: giới hạn bộ nhớ thiết bị tính bằng MiB — công cụ chính để tạo áp lực bộ nhớ.-svà-x: bật chế độ bỏ đầu ra (mô phỏng streaming) và bật profiling bộ nhớ.
Tổng lượng dữ liệu trong cấu hình tiêu chuẩn là 20 GiB mỗi rank, tức 160 GiB trên cả 8 rank.
Shuffle đơn giản: 1,8 TiB/s
Một hệ DGX B200 có 8 GPU Blackwell, mỗi GPU 180 GB VRAM, cùng 2 vi xử lý Intel Xeon Platinum 8570. Khi shuffle dữ liệu nằm gọn trong VRAM, kết quả thật ấn tượng:
[0:PRINT] means: 87.06 ms | local throughput: 229.74 GiB/s | global throughput: 1.79 TiB/s
[7:PRINT] elapsed: 84.58 ms | local throughput: 236.47 GiB/s | global throughput: 1.85 TiB/s
Mỗi rank đạt thông lượng cục bộ khoảng 230 GiB/s, đẩy thông lượng toàn cục lên gần 1,8 TiB/s. Con số này vẫn chưa đạt mức trần lý thuyết 14,4 TB/s, nhưng đã rất nhanh.
Điểm đáng chú ý nằm ở hồ sơ bộ nhớ. Mỗi rank chỉ giữ 20 GiB dữ liệu đầu vào, nhưng mức sử dụng thiết bị đỉnh điểm lên tới 60 GiB — gấp 3 lần đầu vào. Hồ sơ cho thấy rõ: 20 GiB cho chính dữ liệu đầu vào và 40 GiB trong hàm partition_and_pack khi dữ liệu được băm và sao chép vào các buffer đích. Trong một shuffle, chương trình có thể tạm thời sở hữu cả bản gốc lẫn bản sao của dữ liệu cục bộ. Đây là động lực mạnh mẽ để suy nghĩ kỹ về quản lý bộ nhớ ở mọi giai đoạn của pipeline.
Bảng phân bổ bộ nhớ trong quá trình shuffle ngoài bộ nhớ
"Ối, bạn làm tràn mất một chút rồi..."
Để kiểm tra khả năng chịu áp lực bộ nhớ, ta giảm dần giới hạn bộ nhớ thiết bị từ không giới hạn xuống 12 GiB trong khi vẫn giữ 20 GiB dữ liệu mỗi rank. Kết quả cho thấy RapidsMPF không hề gây OOM và cũng không "thrash" — tức là không di chuyển dữ liệu qua lại vô ích giữa thiết bị và host.
Khi giới hạn ở 32 GB, thông lượng toàn cục giảm xuống còn khoảng 480 GiB/s. Điều thú vị nằm ở chỗ: bảng thống kê xuất hiện hai dòng mới là alloc-pinned_host và copy-pinned_host-to-device, nhưng không hề có dòng copy-device-to-pinned_host.
Sự vắng mặt đó chính là điểm mấu chốt. Khi RapidsMPF di chuyển dữ liệu, phía nhận có thể quan sát áp lực bộ nhớ và biết về giới hạn thiết bị đã cấu hình. Thay vì thrash và di chuyển dữ liệu qua lại nhiều lần, phía nhận chấp nhận buffer trên host ngay từ đầu — nhờ UCXX. Nói cách khác, không có gì bị đẩy ra khỏi thiết bị cả; dữ liệu đơn giản là chưa bao giờ hạ cánh xuống đó. Đây là kiểu tràn dữ liệu rẻ nhất: tràn mà không bao giờ phải tràn thật.
Vậy tại sao lại mất hiệu năng nhiều đến vậy? Trong tổng 335,15 ms, RapidsMPF dành tới 171,8 ms chỉ để di chuyển dữ liệu trở lại thiết bị:
copy-pinned_host-to-device: 5.62 GiB | 171.80 ms | 32.74 GiB/s | avg-stream-delay 1.78 ms
DGX B200 dùng PCIe Gen 5 với tốc độ truyền 64 GB/s cho cả 16 làn, và RapidsMPF đang di chuyển dữ liệu qua bus này với tốc độ trung bình 32 GiB/s mỗi GPU. Ngay cả khi dùng bộ nhớ pinned, đây vẫn là nút thắt cổ chai. Nếu thay chip x86 bằng cấu hình Grace-Blackwell với C2C (có thể đạt tới 900 GB/s), ta có thể dễ dàng tăng băng thông giữa host và thiết bị lên 5–10 lần.
Bảng quét mức độ tràn dữ liệu
Khi hạ dần giới hạn bộ nhớ thiết bị, mức độ tràn dữ liệu tăng tương ứng và hiệu năng giảm một cách mượt mà, đơn điệu — nhưng không bao giờ OOM:
- Không tràn (∞ GiB): 229,7 GiB/s cục bộ, 1,79 TiB/s toàn cục.
- Bắt đầu tràn (32 GiB): 60,0 GiB/s cục bộ, 480,2 GiB/s toàn cục.
- Tràn nhẹ (28 GiB): 35,9 GiB/s cục bộ, 287,0 GiB/s toàn cục.
- Tràn vừa (24 GiB): 26,5 GiB/s cục bộ, 211,7 GiB/s toàn cục.
- Tràn nặng (20 GiB): 14,8 GiB/s cục bộ, 118,6 GiB/s toàn cục.
- Tràn rất nặng (16 GiB): 10,1 GiB/s cục bộ, 80,5 GiB/s toàn cục.
- Tràn cực đoan (12 GiB): 6,9 GiB/s cục bộ, 55,5 GiB/s toàn cục.
Ba điểm đáng chú ý rút ra từ bảng này:
- Suy giảm mượt mà và đơn điệu: khi hạ giới hạn bộ nhớ thiết bị, ta dành nhiều thời gian hơn để di chuyển dữ liệu qua lại giữa host và thiết bị ở tốc độ PCIe Gen 5, nhưng điều quan trọng là không hề xảy ra OOM.
- Cơ chế phía nhận tránh được thrash vô ích:
copy-device-to-pinned_hostbằng 0 ở bốn mức tràn đầu tiên. RapidsMPF không bao giờ đẩy dữ liệu ra, mà nhận buffer đến trên host để giảm thiểu giới hạn bộ nhớ thiết bị áp đặt. Chỉ đến giới hạn 20 GiB — khi giới hạn bằng đúng kích thước đầu vào — việc tràn thật từ device sang host mới bắt đầu. - Tràn là phần đắt đỏ nhất:
copy-pinned_host-to-deviceleo từ 5,6 GiB lên 18,8 GiB, và thời gian di chuyển tăng từ 191 ms lên 1170 ms. Với phần cứng có C2C băng thông 900 GB/s, tình hình sẽ khá hơn nhiều.
Kết luận
Shuffle — và những phép join, sort đi kèm — thực sự là bài toán khó và ngốn bộ nhớ. RapidsMPF chứng minh mình có thể xử lý vấn đề này:
- Ngốn bộ nhớ: RapidsMPF có thể tràn/đẩy buffer ra khỏi thiết bị ở các giới hạn có thể tinh chỉnh, đồng thời có thể nhận byte trên host để tránh thrash không cần thiết.
- Truyền tải: RapidsMPF dùng UCXX, cho phép truyền tải qua NVLink, InfiniBand, EFA, TCP, v.v.
- Đồng bộ hóa: RapidsMPF giải quyết thách thức này bằng mô hình thực thi bất đồng bộ/streaming, chồng chéo shuffle với các công việc sẵn sàng khác.
Điều chúng ta hướng tới là một thư viện shuffle nhanh, mở rộng được và xử lý được dữ liệu lớn hơn VRAM, và RapidsMPF đáp ứng được điều đó. Người dùng có thể tái sử dụng thư viện này trực tiếp, hoặc triển khai những ý tưởng này trong dự án của riêng mình. Trong các bài viết tiếp theo, nhóm tác giả hứa hẹn sẽ nghiên cứu khả năng mở rộng trên hệ thống NVL72, cùng với mạng actor và các biểu đồ weak/strong scaling.