IP Library › Granted Patent US 12,229,118
Granted Patent B2
US 12,229,118 · App. 18/742,351 · Granted Feb 18, 2025

Method, apparatus and device for data shuffling, computer-readable storage medium and product

Inventors: Haiyang Shi (Los Angeles, CA); Hao Wang (Los Angeles, CA)
Assignee: Beijing Volcano Engine Technology Co., Ltd.
G06F16/2365G06F7/14
View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 12,229,118
App. No.
18/742,351
Granted
Feb 18, 2025
Kind
B2
Abstract

The embodiments of the disclosure provide a dada shuffling method, apparatus and device, a computer-readable storage medium and product. The method comprises: acquiring a data shuffling request; acquiring a shuffling request parameter linked list associated with the at least one data to be shuffled based on the data shuffling request; performing a merging operation on shuffling request parameters in the shuffling request parameter linked list according to the data amount of the data segment corresponding to the shuffling request parameter and memory buffer information to obtain at least one target request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to a predetermined remote direct memory access network card; and distributing respectively data segments associated with at least one data to be shuffled cached in the remote direct memory access network card to a target server of the data segment.

Claims (74)

1. A method for data shuffling, comprising:

acquiring a data shuffling request for distributing at least one data to be shuffled pre-stored in a host memory to at least one target server, wherein the data to be shuffled comprises at least one data segment to be sent to different target servers;

acquiring a shuffling request parameter linked list associated with the at least one data to be shuffled based on the data shuffling request, wherein the shuffling request parameter linked list comprises at least one shuffling request parameter corresponding to the at least one data segment, and the shuffling request parameter comprises associated information of a target server to which the data segment is to be sent and memory buffer information corresponding to the data segment;

performing a merging operation on shuffling request parameters in the shuffling request parameter linked list according to the data amount of the data segment corresponding to the shuffling request parameter and the memory buffer information to obtain at least one target request parameter composed of at least one shuffling request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to a predetermined remote direct memory access network card; and

distributing respectively data segments associated with at least one data to be shuffled cached in the remote direct memory access network card to the target server indicated in the shuffling request parameter associated with the data segment.

2. The method of claim 1 , wherein the acquiring a shuffling request parameter linked list associated with at least one data to be shuffled in the host memory, comprises:

for each data to be shuffled in the host memory, constructing the shuffling request parameter associated with the data segment based on at least one memory buffer information, the associated information of the target server to send to and a column index corresponding to the data segment; and

generating the shuffling request parameter linked list based on the shuffling request parameter corresponding to each data segment.

3. The method of claim 2 , wherein the shuffling request parameter further comprises a pointer to a next shuffling request parameter;

generating the shuffling request parameter linked list according to the shuffling request parameter corresponding to each data segment, comprises:

based on the pointer carried in each shuffling request parameter, generating the shuffling request parameter linked list according to the shuffling request parameter corresponding to each data segment.

4. The method of claim 1 , wherein performing a merging operation on shuffling request parameters in the shuffling request parameter linked list to obtain at least one target request parameter composed of at least one shuffling request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to the predetermined remote direct memory access network card, comprises:

according to a predetermined request merging condition, merging at least one shuffling request parameter in the shuffling request parameter linked list that meets the request merging condition to obtain at least one target request parameter; and

caching respectively the data to be shuffled corresponding to at least one target request parameter to the remote direct memory access network card by a direct memory access engine predetermined in the remote direct memory access network card.

5. The method of claim 4 , wherein the request merging condition comprises merging the shuffling request parameters whose data amount of the data segment is less than a predetermined data amount threshold and which are continuous on the memory buffer;

merging the shuffling request parameters in the shuffling request parameter linked list according to the data amount of the data segment corresponding to the shuffling request parameter and the memory buffer information, comprises:

for each shuffling request parameter in the shuffling request parameter linked list, detecting whether the data amount of the data segment corresponding to the shuffling request parameter is less than the predetermined data amount threshold; and

merging at least two shuffling request parameters whose data amount is less than the predetermined data amount threshold and which are continuous on the memory buffer to obtain at least one target request parameter.

6. The method of claim 5 , wherein the merging at least two shuffling request parameters whose data amount is less than the predetermined data amount threshold and which are continuous on the memory buffer to obtain at least one target request parameter, further comprises:

for each target request parameter, generating a descriptor of data copy, wherein the descriptor comprises data boundary information and data length information corresponding to each data segment associated with the target request parameter.

7. The method of claim 6 , wherein caching respectively the data to be shuffled corresponding to at least one target request parameter to the remote direct memory access network card by a direct memory access engine predetermined in the remote direct memory access network card, comprises:

copying the data to be shuffled from the host memory according to the target request parameters by the direct memory access engine;

deleting empties in the data to be shuffled based on the descriptor by a predetermined load assembly engine, and assembling data to be sent to the same target server in the target request parameter based on the target server in the target request parameter to obtain assembled data; and

caching respectively the assembled data corresponding to at least one target request parameter to the continuous cache area associated with each target server in the remote direct memory access network card.

8. The method of claim 1 , wherein the shuffling request parameter further comprises predetermined priority information;

caching the data to be shuffled corresponding to the at least one target request parameter to the predetermined remote direct memory access network card, comprises:

determining the direct memory access engine and the load assembly engine corresponding to each shuffling request parameter in the shuffling request parameter linked list based on the priority information;

for each shuffling request parameter, acquiring a target data segment associated with the shuffling request parameter in the host memory by using the direct memory access engine corresponding to the shuffling request parameter; and

caching the target data segment to the continuous cache area predetermined in the predetermined remote direct memory access network card according to the load assembly engine corresponding to the shuffling request parameter.

9. An electronic device, comprising a processor and a memory;

the memory stores computer-executable instructions;

the processor executes computer-executable instructions stored in the memory, to cause the processor to carry out a method for data shuffling, comprising:

acquiring a data shuffling request for distributing at least one data to be shuffled pre-stored in a host memory to at least one target server, wherein the data to be shuffled comprises at least one data segment to be sent to different target servers;

acquiring a shuffling request parameter linked list associated with the at least one data to be shuffled based on the data shuffling request, wherein the shuffling request parameter linked list comprises at least one shuffling request parameter corresponding to the at least one data segment, and the shuffling request parameter comprises associated information of a target server to which the data segment is to be sent and memory buffer information corresponding to the data segment;

performing a merging operation on shuffling request parameters in the shuffling request parameter linked list according to the data amount of the data segment corresponding to the shuffling request parameter and the memory buffer information to obtain at least one target request parameter composed of at least one shuffling request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to a predetermined remote direct memory access network card; and

distributing respectively data segments associated with at least one data to be shuffled cached in the remote direct memory access network card to the target server indicated in the shuffling request parameter associated with the data segment.

10. The electronic device of claim 9 , wherein the acquiring a shuffling request parameter linked list associated with at least one data to be shuffled in the host memory, comprises:

for each data to be shuffled in the host memory, constructing the shuffling request parameter associated with the data segment based on at least one memory buffer information, the associated information of the target server to send to and a column index corresponding to the data segment; and

generating the shuffling request parameter linked list based on the shuffling request parameter corresponding to each data segment.

11. The electronic device of claim 10 , wherein the shuffling request parameter further comprises a pointer to a next shuffling request parameter;

generating the shuffling request parameter linked list according to the shuffling request parameter corresponding to each data segment, comprises:

based on the pointer carried in each shuffling request parameter, generating the shuffling request parameter linked list according to the shuffling request parameter corresponding to each data segment.

12. The electronic device of claim 9 , wherein performing a merging operation on shuffling request parameters in the shuffling request parameter linked list obtain at least one target request parameter composed of at least one shuffling request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to the predetermined remote direct memory access network card, comprises:

according to a predetermined request merging condition, merging at least one shuffling request parameter in the shuffling request parameter linked list that meets the request merging condition to obtain at least one target request parameter; and

caching respectively the data to be shuffled corresponding to at least one target request parameter to the remote direct memory access network card by a direct memory access engine predetermined in the remote direct memory access network card.

13. The electronic device of claim 12 , wherein the request merging condition comprises merging the shuffling request parameters whose data amount of the data segment is less than a predetermined data amount threshold and which are continuous on the memory buffer;

merging the shuffling request parameters in the shuffling request parameter linked list according to the data amount of the data segment corresponding to the shuffling request parameter and the memory buffer information, comprises:

for each shuffling request parameter in the shuffling request parameter linked list, detecting whether the data amount of the data segment corresponding to the shuffling request parameter is less than the predetermined data amount threshold; and

merging at least two shuffling request parameters whose data amount is less than the predetermined data amount threshold and which are continuous on the memory buffer to obtain at least one target request parameter.

14. The electronic device of claim 13 , wherein the merging at least two shuffling request parameters whose data amount is less than the predetermined data amount threshold and which are continuous on the memory buffer to obtain at least one target request parameter, further comprises:

for each target request parameter, generating a descriptor of data copy, wherein the descriptor comprises data boundary information and data length information corresponding to each data segment associated with the target request parameter.

15. The electronic device of claim 14 , wherein caching respectively the data to be shuffled corresponding to at least one target request parameter to the remote direct memory access network card by a direct memory access engine predetermined in the remote direct memory access network card, comprises:

copying the data to be shuffled from the host memory according to the target request parameters by the direct memory access engine;

deleting empties in the data to be shuffled based on the descriptor by a predetermined load assembly engine, and assembling data to be sent to the same target server in the target request parameter based on the target server in the target request parameter to obtain assembled data; and

caching respectively the assembled data corresponding to at least one target request parameter to the continuous cache area associated with each target server in the remote direct memory access network card.

16. The electronic device of claim 9 , wherein the shuffling request parameter further comprises predetermined priority information;

caching the data to be shuffled corresponding to the at least one target request parameter to the predetermined remote direct memory access network card, comprises:

determining the direct memory access engine and the load assembly engine corresponding to each shuffling request parameter in the shuffling request parameter linked list based on the priority information;

for each shuffling request parameter, acquiring a target data segment associated with the shuffling request parameter in the host memory by using the direct memory access engine corresponding to the shuffling request parameter; and

caching the target data segment to the continuous cache area predetermined in the predetermined remote direct memory access network card according to the load assembly engine corresponding to the shuffling request parameter.

17. A non-transitory computer-readable storage medium, wherein the computer-readable storage medium stores computer-executable instructions which, when executed by a processor, carry out a method for data shuffling comprising:

acquiring a data shuffling request for distributing at least one data to be shuffled pre-stored in a host memory to at least one target server, wherein the data to be shuffled comprises at least one data segment to be sent to different target servers;

acquiring a shuffling request parameter linked list associated with the at least one data to be shuffled based on the data shuffling request, wherein the shuffling request parameter linked list comprises at least one shuffling request parameter corresponding to the at least one data segment, and the shuffling request parameter comprises associated information of a target server to which the data segment is to be sent and memory buffer information corresponding to the data segment;

performing a merging operation on shuffling request parameters in the shuffling request parameter linked list according to the data amount of the data segment corresponding to the shuffling request parameter and the memory buffer information to obtain at least one target request parameter composed of at least one shuffling request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to a predetermined remote direct memory access network card; and

distributing respectively data segments associated with at least one data to be shuffled cached in the remote direct memory access network card to the target server indicated in the shuffling request parameter associated with the data segment.

18. The non-transitory computer-readable storage medium of claim 17 , wherein the acquiring a shuffling request parameter linked list associated with at least one data to be shuffled in the host memory, comprises:

for each data to be shuffled in the host memory, constructing the shuffling request parameter associated with the data segment based on at least one memory buffer information, the associated information of the target server to send to and a column index corresponding to the data segment; and

generating the shuffling request parameter linked list based on the shuffling request parameter corresponding to each data segment.

19. The non-transitory computer-readable storage medium of claim 18 , wherein the shuffling request parameter further comprises a pointer to a next shuffling request parameter;

generating the shuffling request parameter linked list according to the shuffling request parameter corresponding to each data segment, comprises:

based on the pointer carried in each shuffling request parameter, generating the shuffling request parameter linked list according to the shuffling request parameter corresponding to each data segment.

20. The non-transitory computer-readable storage medium of claim 17 , wherein performing a merging operation on shuffling request parameters in the shuffling request parameter linked list to obtain at least one target request parameter composed of at least one shuffling request parameter; and caching the data to be shuffled corresponding to the at least one target request parameter to the predetermined remote direct memory access network card, comprises:

according to a predetermined request merging condition, merging at least one shuffling request parameter in the shuffling request parameter linked list that meets the request merging condition to obtain at least one target request parameter; and

caching respectively the data to be shuffled corresponding to at least one target request parameter to the remote direct memory access network card by a direct memory access engine predetermined in the remote direct memory access network card.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 3, 2024
From: SHI, HAIYANG; WANG, HAO
To: BYTEDANCE INC.
Reel/Frame 069472/0395 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 3, 2024
From: BYTEDANCE INC.
To: BEIJING VOLCANO ENGINE TECHNOLOGY CO., LTD.
Reel/Frame 069472/0544 →
Priority Claims (1)
CN 202310879264.8 · Jul 17, 2023 · national
Continuity (1)
Related Publication 20250028705A1 · Jan 23, 2025
References Cited (17)
US 9256536B2 · Park · 2016 [cited by examiner]
US 20060045108A1 · Blackmore · 2006 [cited by examiner]
US 20060045109A1 · Blackmore · 2006 [cited by examiner]
US 20060075067A1 · Blackmore · 2006 [cited by examiner]
US 20110313973A1 · Srivas · 2011 [cited by examiner]
US 20150150018A1 · Hu · 2015 [cited by examiner]
US 20150281126A1 · Regula et al. · 2015 [cited by applicant]
US 20200133533A1 · Zhao · 2020 [cited by examiner]
US 20200202197A1 · Subhaschandra Banakar · 2020 [cited by examiner]
US 20200341764A1 · Jacob et al. · 2020 [cited by applicant]
US 20220164122A1 · Zou · 2022 [cited by examiner]
US 20230125593A1 · Mahony · 2023 [cited by examiner]
US 20230244629A1 · Marcovitch · 2023 [cited by examiner]
CN 103647807A · 2014 [cited by applicant]
CN 103902486A · 2017 [cited by applicant]
B. Liu, F. Liu, N. Xiao and Z. Chen, “Accelerating Spark Shuffle with RDMA,” 2018 IEEE International Conference on Networking, Architecture and Storage (NAS), Chongqing, China, 2018, pp. 1-7. (Year: 2018). [cited by examiner]
European Patent Office, Extended European Search Report Issued in Application No. 24181560.4, Nov. 28, 2024, 10 pages. [cited by applicant]