IP Library Granted Patent US 10,642,520
Granted Patent B1
US 10,642,520 · App. 15/490,430 · Granted May 5, 2020

Memory optimized data shuffle

Inventors: Junping Zhao (Beijing, CN); Kenneth J. Taylor (Franklin, MA); Randall Shain (Wrentham, MA); Kun Wang (Beijing, CN)
Assignee: EMC IP Holding Company LLC
G06F3/0638G06F3/061G06F3/065G06F3/067G06F3/0647G06F3/0655G06F3/0656G06F3/0683
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 10,642,520
App. No.
15/490,430
Granted
May 5, 2020
Kind
B1
Abstract

In a distributed data processing system with a set of multiple nodes, a first data shuffle memory pool is maintained at a data shuffle writer node, and a second data shuffle memory pool is maintained at a data shuffle reader node. The data shuffle writer node and the data shuffle reader node are part of the set of multiple nodes of the distributed data processing system. In-memory compression is performed on at least a portion of a data set from the first data shuffle memory pool. At least a portion of the compressed data is transmitted from the first data shuffle memory pool to the second data shuffle memory pool in a peer-to-peer manner. Each of the first data shuffle memory pool and the second data shuffle memory pool may include a hybrid memory configuration.

Claims (40)

1. A method comprising:

maintaining a first data shuffle memory pool at a data shuffle writer node and a second data shuffle memory pool at a data shuffle reader node, wherein the data shuffle writer node and the data shuffle reader node are part of a set of multiple nodes of a distributed data processing system;

performing an in-memory compression on at least a portion of a data set from the first data shuffle memory pool;

performing a data shuffle operation on the at least a portion of the compressed data from the first shuffle memory pool, wherein the data shuffle operation maps different parts of the compressed data for transmission to different nodes of the distributed data processing system; and

transmitting, in response to the data shuffle operation, the at least a portion of the compressed data from the first data shuffle memory pool to the second data shuffle memory pool in a peer-to-peer manner;

wherein the distributed data processing system is implemented via one or more processing devices.

2. The method of claim 1 , wherein the performance of the in-memory compression is triggered by a given policy.

3. The method of claim 2 , wherein the given policy is based on a memory pressure associated with a memory capacity threshold.

4. The method of claim 2 , wherein the given policy is based on a ratio of incoming write speed to outgoing read speed.

5. The method of claim 1 , wherein performing the in-memory compression further comprises:

obtaining one or more records;

compressing the obtained records into a work buffer;

copying the compressed records to an original location of the one or more compressed records; and

updating an active location, the active location indicating the beginning location of a next non-compressed record.

6. The method of claim 1 , wherein each of the first data shuffle memory pool and the second data shuffle memory pool comprise a hybrid memory configuration.

7. The method of claim 6 , wherein the hybrid memory configuration comprises a dynamic random access memory tier and one or more nonvolatile random access memory tiers.

8. The method of claim 7 , wherein the one or more nonvolatile random access memory tiers comprise a local nonvolatile random access memory tier and a remote nonvolatile random access memory tier.

9. The method of claim 7 , wherein the one or more nonvolatile random access memory tiers comprise a global nonvolatile random access memory tier.

10. The method of claim 1 , wherein each of the first data shuffle memory pool and the second data shuffle memory pool comprise a global shared memory tier.

11. The method of claim 1 , wherein data in the data set is characterized by the data shuffle writer node as one or more of: data to be transmitted to the data shuffle reader node; data to be infrequently locally accessed; and data to be sorted and aggregated as one or more extents.

12. The method of claim 1 , wherein the portion of the data set transmitted from the first data shuffle memory pool to the second data shuffle memory pool is transferred using remote direct memory access.

13. The method of claim 1 , wherein the portion of the data set to be transmitted is sorted in the data shuffle writer node based on one or more memory ranges associated with the data shuffle reader node.

14. The method of claim 1 , wherein the distributed data processing system comprises a MapReduce processing framework, and wherein the data shuffle writer node comprises a mapper node and the data shuffle read node comprises a reducer node.

15. A distributed data processing system, comprising:

a plurality of nodes wherein at least one node is configured as a data shuffle writer node and at least another node is configured as a data shuffle reader node, wherein:

the data shuffle writer node is configured to maintain a first data shuffle memory pool;

the data shuffle reader node is configured to maintain a second data shuffle memory pool; and further wherein:

the data shuffle writer node is configured to perform an in-memory compression on at least a portion of a data set from the first data shuffle memory pool;

the data shuffle writer node is configured to perform a data shuffle operation on at least a portion of the compressed data from the first shuffle memory pool, wherein the data shuffle operation maps different parts of the compressed data for transmission to different nodes of the distributed data processing system; and

the data shuffle writer node, in response to the data shuffle operation, is configured to transmit the at least a portion of the compressed data from the first data shuffle memory pool to the second data shuffle memory pool in a peer-to-peer manner;

wherein the distributed data processing system is implemented via one or more processing devices.

16. The system of claim 15 , wherein performance of the in-memory compression is triggered by a given policy, and wherein the given policy is based on a memory pressure associated with a memory capacity threshold or a ratio of incoming write speed to outgoing read speed.

17. The system of claim 15 , wherein each of the first data shuffle memory pool and the second data shuffle memory pool comprise a hybrid memory configuration.

18. The system of claim 15 , wherein the portion of the data set transmitted from the first data shuffle memory pool to the second data shuffle memory pool is transferred using remote direct memory access.

19. The system of claim 15 , wherein the distributed data processing system comprises a MapReduce processing framework, and wherein the data shuffle writer node comprises a mapper node and the data shuffle read node comprises a reducer node.

20. A computer program product comprising a non-transitory processor-readable storage medium having stored therein program code of one or more software programs, wherein the program code when executed by one or more processing devices of a distributed data processing system causes said one or more processing devices to:

maintain a first data shuffle memory pool at a data shuffle writer node and a second data shuffle memory pool at a data shuffle reader node, wherein the data shuffle writer node and the data shuffle reader node are part of a set of multiple nodes of the distributed data processing system;

perform an in-memory compression on at least a portion of a data set from the first data shuffle memory pool; and

perform a data shuffle operation on the at least a portion of the compressed data from the first shuffle memory pool, wherein the data shuffle operation maps different parts of the compressed data for transmission to different nodes of the distributed data processing system; and

transmit, in response to the data shuffle operation, the at least a portion of the compressed data from the first data shuffle memory pool to the second data shuffle memory pool in a peer-to-peer manner.

Assignments (8)
RELEASE OF SECURITY INTEREST IN PATENTS PREVIOUSLY RECORDED AT REEL/FRAME (053546/0001) Recorded Jun 23, 2022
From: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS NOTES COLLATERAL AGENT
To: DELL MARKETING L.P. (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO CREDANT TECHNOLOGIES, INC.); DELL INTERNATIONAL L.L.C.; DELL PRODUCTS L.P.; DELL USA L.P.; EMC CORPORATION; DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO FORCE10 NETWORKS, INC. AND WYSE TECHNOLOGY L.L.C.); EMC IP HOLDING COMPANY LLC
Reel/Frame 071642/0001 →
RELEASE OF SECURITY INTEREST IN PATENTS PREVIOUSLY RECORDED AT REEL/FRAME (042769/0001) Recorded Apr 26, 2022
From: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS NOTES COLLATERAL AGENT
To: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO MOZY, INC.); DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO WYSE TECHNOLOGY L.L.C.)
Reel/Frame 059803/0802 →
RELEASE OF SECURITY INTEREST AT REEL 042768 FRAME 0585 Recorded Nov 2, 2021
From: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH
To: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; MOZY, INC.; WYSE TECHNOLOGY L.L.C.
Reel/Frame 058297/0536 →
SECURITY AGREEMENT Recorded Apr 22, 2020
From: CREDANT TECHNOLOGIES INC.; DELL INTERNATIONAL L.L.C.; DELL MARKETING L.P.; DELL PRODUCTS L.P.; DELL USA L.P.; EMC CORPORATION; FORCE10 NETWORKS, INC.; WYSE TECHNOLOGY L.L.C.; EMC IP HOLDING COMPANY LLC
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A.
Reel/Frame 053546/0001 →
SECURITY AGREEMENT Recorded Mar 21, 2019
From: CREDANT TECHNOLOGIES, INC.; DELL INTERNATIONAL L.L.C.; DELL MARKETING L.P.; DELL PRODUCTS L.P.; DELL USA L.P.; EMC CORPORATION; FORCE10 NETWORKS, INC.; WYSE TECHNOLOGY L.L.C.; EMC IP HOLDING COMPANY LLC
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A.
Reel/Frame 049452/0223 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jul 24, 2017
From: ZHAO, JUNPING; TAYLOR, KENNETH J.; SHAIN, RANDALL; WANG, KUN
To: EMC IP HOLDING COMPANY LLC
Reel/Frame 043081/0283 →
PATENT SECURITY INTEREST (CREDIT) Recorded Jun 12, 2017
From: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; MOZY, INC.; WYSE TECHNOLOGY L.L.C.
To: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH, AS COLLATERAL AGENT
Reel/Frame 042768/0585 →
PATENT SECURITY INTEREST (NOTES) Recorded Jun 12, 2017
From: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; MOZY, INC.; WYSE TECHNOLOGY L.L.C.
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS COLLATERAL AGENT
Reel/Frame 042769/0001 →
Cited By (2)
US 12,229,168 US 12,468,443