IP Library Granted Patent US 10,698,891
Granted Patent B2
US 10,698,891 · App. 15/668,861 · Granted Jun 30, 2020

MxN dispatching in large scale distributed system

Inventors: Lei Chang (Beijing, CN); Tao Ma (Beijing, CN); Zhanwei Wang (Beijing, CN); Luke Lonergan (San Carlos, CA); Lirong Jian (Beijing, CN); Lili Ma (Beijing, CN)
Assignee: EMC IP Holding Company LLC
G06F16/24542G06F16/10G06F16/11G06F16/148G06F16/182G06F16/1858G06F16/2453G06F16/2455G06F16/2471G06F16/24524G06F16/24532G06F16/27G06F16/907H04L65/60H04L67/1097H05K999/99G06F16/113G06F16/217G06F16/245G06F16/43
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,698,891
App. No.
15/668,861
Granted
Jun 30, 2020
Kind
B2
Abstract

M×N dispatching in a large scale distributed system is disclosed. In various embodiments, a query is received. A query plan is generated to perform the query. A subset of query processing segments is selected, from a set of available query processing segments, to perform an assigned portion of the query plan. An assignment to perform the assigned portion of the query plan is dispatched to the selected subset of query processing segments.

Claims (34)

1. A method, comprising:

receiving a query;

generating, by a master node, a query plan to perform the query, wherein the generating of the query plan includes dividing the query plan into at least a first portion and a second portion, and wherein the master node comprises one or more hardware processors;

selecting, by the master node, from a set of available query processing segments a first subset of query processing segments to perform a first assigned portion of the query plan corresponding to the first portion of the query plan, and a second subset of query processing segments to perform a second assigned portion of the query plan corresponding to the second portion of the query plan; and

dispatching to the selected first subset of query processing segments an assignment to perform the first assigned portion of the query plan, wherein the dispatching of the assignment to perform the first assigned portion of the query plan includes providing to the selected first subset of query processing segments with corresponding metadata that is obtained from a central metadata store, wherein the metadata provided to the corresponding selected first subset of query processing segments is determined to be used by the selected first subset of query processing segments to perform the first assigned portion of the query plan.

2. The method of claim 1 , wherein a first number of segments selected to perform the first portion of the query plan is dynamically determined according to one or both of (1) data locality of data corresponding to the first portion of the query plan associated with the query in relation to the first subset of query processing segments, and (2) available resources.

3. The method of claim 2 , wherein the first subset of query processing segments is selected based at least in part on a co-locality of one or more of the selected query processing segments with data with which the assigned portion of the query plan is associated.

4. The method of claim 2 , wherein the first number of segments is selected to perform the first portion of the query plan and a second number of segments, different from the first number, is selected to perform the second portion of the query plan.

5. The method of claim 1 , wherein the first subset of query processing segments and the second subset of query processing segments include at least a first segment.

6. The method of claim 1 , wherein the metadata to be used by the selected first subset of query processing segments to perform the first assigned portion of the query plan is embedded in a communication from the master node to the selected first subset of query processing segments that comprises the first assigned portion of the query plan.

7. The method of claim 1 , wherein the metadata comprises data indicating where data required to perform first assigned portion of the query plan is located.

8. The method of claim 1 , wherein a first number of segments selected to perform the first portion of the query plan is dynamically determined.

9. The method of claim 1 , wherein the first subset of query processing segments executes a plurality of query execution threads.

10. The method of claim 1 , wherein the dividing the query plan into at least a first portion and a second portion includes dividing the query plan into a plurality of independently executable slices.

11. The method of claim 1 , wherein selecting the first subset of query processing segments includes receiving from a resource manager an indication of a degree of availability of processing segments included in the set of available query processing segments.

12. The method of claim 1 , wherein the first subset of query processing segments is selected according to a determination for each of a plurality of portions of the query plan a corresponding number of segments to be assigned to perform that portion.

13. The method of claim 1 , wherein the assignment comprises a network communication sent via a network interconnect.

14. The method of claim 1 , wherein the assignment includes the metadata, wherein the metadata is embedded in the assignment and is to be used to perform one or more tasks associated with the assigned portion of the query plan.

15. The method of claim 1 , wherein segments comprising the subset of query processing segments are associated with one or more segments hosts, each of which is configured to provide one or more processing segments.

16. The method of claim 1 , wherein the assignment includes the metadata, wherein the metadata indicates a location, within a distributed storage layer, of data associated with the assignment.

17. The method of claim 16 , wherein the location includes identification of a table of data to be search in connection with the assignment.

18. The method of claim 1 , wherein the first subset of query processing segments comprises one or more query processing segments, and the second subset of query processing segments comprises one or more query processing segments.

19. A system, comprising:

a communication interface; and

one or more hardware processors coupled to the communication interface and configured to:

receive a query;

generate a query plan to perform the query, wherein the query plan is generated such that the query plan is divided into at least a first portion and a second portion;

select from a set of available query processing segments a first subset of query processing segments to perform a first assigned portion of the query plan, corresponding to the first portion of the query plan, and a second subset of query processing segments to perform a second assigned portion of the query plan corresponding to the second portion of the query plan; and

dispatch to the selected first subset of query processing segments, via the communication interface, an assignment to perform the first assigned portion of the query plan, wherein to dispatch the assignment to perform the first assigned portion of the query plan includes providing to the selected first subset of query processing segments with corresponding metadata that is obtained from a central metadata store, wherein the metadata provided to the corresponding selected first subset of query processing segments is determined to be used by the selected first subset of query processing segments to perform the first assigned portion of the query plan.

20. A computer program product embodied in a tangible, non-transitory computer readable storage means, comprising computer instructions for:

receiving a query;

generating a query plan to perform the query, wherein the generating of the query plan includes dividing the query plan into at least a first portion and a second portion;

selecting from a set of available query processing segments a first subset of query processing segments to perform a first assigned portion of the query plan corresponding to the first portion of the query plan, and a second subset of query processing segments to perform a second assigned portion of the query plan corresponding to the second portion of the query plan; and

dispatching to the selected first subset of query processing segments an assignment to perform the first assigned portion of the query plan, wherein the dispatching of the assignment to perform the first assigned portion of the query plan includes providing to the selected first subset of query processing segments with corresponding metadata that is obtained from a central metadata store, wherein the metadata provided to the corresponding selected first subset of query processing segments is determined to be used by the selected first subset of query processing segments to perform the first assigned portion of the query plan.

Assignments (9)
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 (045482/0131) Recorded May 20, 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; DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO WYSE TECHNOLOGY L.L.C.)
Reel/Frame 061749/0924 →
RELEASE OF SECURITY INTEREST AT REEL 045482 FRAME 0395 Recorded Nov 2, 2021
From: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH
To: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; WYSE TECHNOLOGY L.L.C.
Reel/Frame 058298/0314 →
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 →
PATENT SECURITY AGREEMENT (CREDIT) Recorded Mar 1, 2018
From: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; WYSE TECHNOLOGY L.L.C.
To: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH, AS COLLATERAL AGENT
Reel/Frame 045482/0395 →
PATENT SECURITY AGREEMENT (NOTES) Recorded Mar 1, 2018
From: DELL PRODUCTS L.P.; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; WYSE TECHNOLOGY L.L.C.
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS COLLATERAL AGENT
Reel/Frame 045482/0131 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 16, 2017
From: CHANG, LEI; MA, TAO; WANG, ZHANWEI; LONERGAN, LUKE; JIAN, LIRONG; MA, LILI
To: EMC CORPORATION
Reel/Frame 043872/0452 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 16, 2017
From: EMC CORPORATION
To: EMC IP HOLDING COMPANY LLC
Reel/Frame 044212/0572 →