IP Library Granted Patent US 11,023,274
Granted Patent B2
US 11,023,274 · App. 16/711,672 · Granted Jun 1, 2021

Method and system for processing data

Inventors: Simon Tao (Shanghai, CN); Yu Cao (Beijing, CN); Zhe Dong (Beijing, CN); Sanping Li (Beijing, CN)
Assignee: EMC IP Holding Company LLC
G06F9/4856G06F9/4881G06F9/5083G06F16/9024G06K9/6296G06K9/6892
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 11,023,274
App. No.
16/711,672
Granted
Jun 1, 2021
Kind
B2
Abstract

A method for processing data includes receiving an adjustment request for adjusting a number of consumer instances from a first number to a second number, and determining a migration overhead for adjusting a first distribution of states associated with the first number of consumer instances to a second distribution of the states associated with the second number of consumer instances, wherein the states are intermediate results of processing the data and the migration overhead includes a latency and a bandwidth shortage incurred for migrating the states. Based on the determined migration overhead, the states are migrated between the first number of consumer instances and the second number of consumer instances, and thereafter the data is processed based on the second distribution of the states at the second number of consumer instances.

Claims (42)

1. A method for processing data, comprising:

receiving an adjustment request for adjusting a number of consumer instances from a first number of consumer instances to a second number of consumer instances;

determining a migration overhead for adjusting a first distribution of states associated with the first number of consumer instances to a second distribution of the states associated with the second number of consumer instances, the states being intermediate results of processing the data, the determining being based on a consistent hashing algorithm identifying K/n states to be migrated, wherein K is a total number of states across the first number of consumer instances and n is the second number;

based on the determined migration overhead, migrating the states between the first number of the consumer instances and the second number of the consumer instances; and

thereafter processing the data based on the second distribution of the states at the second number of the consumer instances.

2. The method according to claim 1 , further comprising:

at one consumer instance of the second number of consumer instances before the migrating of the states is completed, caching data dispatched from a producer instance and processing the cached data after a relevant portion of the states are migrated.

3. The method according to claim 2 , wherein the data is dispatched to the second number of consumer instances according to the second distribution of the states.

4. The method according to claim 2 , wherein the relevant portion of the states are migrated to the one consumer instance in response to receipt of a set of adjustment indicators for the adjustment of the states among the consumer instances.

5. The method according to claim 4 , further comprising:

before the set of adjustment indicators are received, processing the data at the one consumer instance based on the first distribution of the states.

6. The method according to claim 1 , further including at least one of:

in response to determining that a workload for processing the data increases, sending the adjustment request for increasing the number of the consumer instances from the first number to the second number; or

in response to determining that a workload for processing the data decreases, sending the adjustment request for decreasing the number of the consumer instances from the first number to the second number.

7. The method according to claim 1 , further comprising:

dynamically adjusting computing resources at runtime of a stream processing application which operates in accordance with a Directed Acyclic Graph (DAG) topology.

8. A system for processing data, comprising:

one or more processors;

a memory coupled to at least one processor of the one or more processors, the memory storing computer program instructions which, when executed by the processors, cause the system to execute a method for processing data, the method comprising:

receiving an adjustment request for adjusting a number of consumer instances from a first number to a second number;

determining a migration overhead for adjusting a first distribution of states associated with the first number of consumer instances to a second distribution of the states associated with the second number of consumer instances, the states being intermediate results of processing the data, the determining being based on a consistent hashing algorithm identifying K/n states to be migrated, wherein K is a total number of states across the first number of consumer instances and n is the second number;

based on the determined migration overhead, migrating the states between the first number of the consumer instances and the second number of the consumer instances; and

thereafter processing the data based on the second distribution of the states at the second number of the consumer instances.

9. The system according to claim 8 , further comprising:

at one consumer instance of the second number of consumer instances before the migrating of the states is completed, caching data dispatched from a producer instance and processing the cached data after a relevant portion of the states are migrated.

10. The system according to claim 9 , wherein the data is dispatched to the second number of consumer instances according to the second distribution of the states.

11. The system according to claim 9 , wherein the relevant portion of the states are migrated to the one consumer instance in response to receipt of a set of adjustment indicators for the adjustment of the states among the consumer instances.

12. The system according to claim 11 , wherein the method further includes:

before the set of adjustment indicators are received, processing the data at the one consumer instance based on the first distribution of the states.

13. The system according to claim 8 , wherein the method further includes at least one of:

in response to determining that a workload for processing the data increases, sending the adjustment request for increasing the number of the consumer instances from the first number to the second number; or

in response to determining that a workload for processing the data decreases, sending the adjustment request for decreasing the number of the consumer instances from the first number to the second number.

14. The system according to claim 8 , wherein the method further includes:

dynamically adjusting computing resources at runtime of a stream processing application which operates in accordance with a Directed Acyclic Graph (DAG) topology.

15. A computer program product having a non-transitory computer readable medium which stores a set of instructions which, when carried out by computerized circuitry, cause the computerized circuitry to perform a method of processing data, the method comprising:

receiving an adjustment request for adjusting a number of consumer instances from a first number to a second number;

determining a migration overhead for adjusting a first distribution of states associated with the first number of consumer instances to a second distribution of the states associated with the second number of consumer instances, the states being intermediate results of processing the data, the determining being based on a consistent hashing algorithm identifying K/n states to be migrated, wherein K is a total number of states across the first number of consumer instances and n is the second number;

based on the determined migration overhead, migrating the states between the first number of the consumer instances and the second number of the consumer instances; and

thereafter processing the data based on the second distribution of the states at the second number of the consumer instances.

16. The computer program product according to claim 15 , wherein the method further includes:

at one consumer instance of the second number of consumer instances before the migrating of the states is completed, caching data dispatched from a producer instance and processing the cached data after a relevant portion of the states are migrated.

17. The computer program product according to claim 16 , wherein the data is dispatched to the second number of consumer instances according to the second distribution of the states.

Assignments (11)
RELEASE OF SECURITY INTEREST IN PATENTS PREVIOUSLY RECORDED AT REEL/FRAME (053311/0169) Recorded Jun 23, 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
Reel/Frame 060438/0742 →
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 (052216/0758) Recorded Jun 23, 2022
From: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS NOTES COLLATERAL AGENT
To: DELL PRODUCTS L.P.; EMC IP HOLDING COMPANY LLC
Reel/Frame 060438/0680 →
RELEASE OF SECURITY INTEREST AF REEL 052243 FRAME 0773 Recorded Nov 2, 2021
From: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH
To: DELL PRODUCTS L.P.; EMC IP HOLDING COMPANY LLC
Reel/Frame 058001/0152 →
SECURITY INTEREST Recorded Jun 5, 2020
From: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS COLLATERAL AGENT
Reel/Frame 053311/0169 →
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 26, 2020
From: DELL PRODUCTS L.P.; EMC IP HOLDING COMPANY LLC
To: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH
Reel/Frame 052243/0773 →
PATENT SECURITY AGREEMENT (NOTES) Recorded Mar 24, 2020
From: DELL PRODUCTS L.P.; EMC IP HOLDING COMPANY LLC
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS COLLATERAL AGENT
Reel/Frame 052216/0758 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 30, 2020
From: EMC CORPORATION
To: EMC IP HOLDING COMPANY LLC
Reel/Frame 051757/0001 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 24, 2020
From: DONG, ZHE
To: EMC CORPORATION
Reel/Frame 051605/0198 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 24, 2020
From: TAO, SIMON; CAO, YU; LI, SANPING
To: EMC IP HOLDING COMPANY LLC
Reel/Frame 051605/0202 →