IP Library › Granted Patent US 10,015,106
Granted Patent B1
US 10,015,106 · App. 14/982,341 · Granted Jul 3, 2018

Multi-cluster distributed data processing platform

Inventors: Patricia Gomes Soares Florissi (Briarcliff Manor, NY); Benny Lutati (Beer-Sheva, IL); Ehud Gudes (Beer-Sheva, IL); Yaron Gonen (Rehovot, IL); Ido Singer (Nes-Ziona, IL); Amnon Meisels (Omer, IL); Sudhir Vijendra (Westborough, MA)
Assignee: EMC IP Holding Company LLC
H04L47/70H04L67/10
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,015,106
App. No.
14/982,341
Filed
Dec 29, 2015
Granted
Jul 3, 2018
Kind
B1
Art Unit
2446
USPC
709/226
Abstract

An apparatus comprises a multi-cluster distributed data processing platform. Each of the clusters of the platform is configured to perform processing operations utilizing local data resources accessible within a corresponding data zone. A first application is initiated in one of the clusters, and data resources to be utilized by the application are determined. For one or more of the data resources identified as local data resources for the associated cluster, processing operations are performed utilizing those local data resources in that cluster in accordance with the first application. For one or more of the data resources identified as remote data resources for the associated cluster, one or more additional applications are initiated in one or more additional ones of the clusters. This repeats recursively for each additional application until all processing required by the first application is complete. Processing results from the clusters are aggregated and provided to a client.

Claims (73)

1. A method comprising:

initiating a first application in a first one of a plurality of distributed processing node clusters associated with respective data zones, each of the clusters being configured to perform processing operations utilizing local data resources locally accessible within its corresponding data zone;

determining a plurality of data resources to be utilized by the application;

identifying for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource that is locally accessible within the data zone of the first distributed processing node cluster or a remote data resource that is not locally accessible within the data zone of the first distributed processing node cluster;

for one or more of the plurality of data resources that are identified as local data resources, performing processing operations utilizing the local data resources in the first cluster in accordance with the first application;

for one or more of the plurality of data resources that are identified as remote data resources, initiating respective additional applications in one or more additional ones of the plurality of distributed processing node clusters and performing processing operations utilizing the remote data resources in the corresponding one or more additional clusters in accordance with the one or more additional applications;

aggregating processing results from the first and one or more additional clusters; and

providing the aggregated processing results to a client;

wherein the processing results from the first and one or more additional clusters preserve at least one specified policy of those clusters in the respective local and remote data resources wherein the specified policy relates to at least one of privacy, security, governance, risk and compliance;

wherein the processing results from a given one of the clusters are permitted to be transmitted to another one of the clusters but the local data resources of the given cluster that are utilized to obtain the processing results are not permitted to be transmitted to another one of the clusters; and

wherein the method is implemented by at least one processing device comprising a processor coupled to a memory.

2. The method of claim 1 wherein the plurality of distributed processing node clusters comprise respective YARN clusters.

3. The method of claim 1 wherein the clusters comprise respective geographically-distributed regional data centers each configured to perform analytics processing utilizing the locally accessible data resources of its corresponding data zone.

4. The method of claim 1 wherein the method is performed at least in part in a worldwide data node coupled to one or more of the clusters.

5. The method of claim 1 wherein the method is performed at least in part in a worldwide data node that comprises a processing node of a given one of the clusters.

6. The method of claim 1 wherein the client initiates the first application in the first one of the clusters as a YARN application and the first cluster initiates the one or more additional applications in the one or more additional clusters as respective YARN applications for which the first cluster serves as a client such that the one or more additional clusters are unaware that the one or more additional applications are part of a multi-cluster distributed application.

7. The method of claim 1 further comprising accessing a distributed catalog service to identify for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource or a remote data resource.

8. The method of claim 1 wherein at least one of the additional clusters determines an additional plurality of data resources to be utilized by the corresponding additional application and identifies for each of the plurality of additional data resources to be utilized by the additional application whether the data resource is a local data resource that is locally accessible within the data zone of the additional cluster or a remote data resource that is not locally accessible within the data zone of the additional cluster.

9. The method of claim 8 wherein the additional plurality of data resources includes one or more remote data resources not locally accessible to the additional cluster and the additional cluster initiates one or more other applications in one or more other ones of the clusters that have local access to the one or more remote data resources.

10. The method of claim 1 wherein a given one of the applications initiated in a corresponding one of the clusters comprises:

an application master component;

at least one aggregator controlled by the application master component; and

one or more containers utilized by the aggregator in conjunction with execution of the given application;

wherein the one or more containers correspond to respective ones of the local data resources that are locally accessible within the corresponding cluster; and

wherein the aggregator is configured to request initiation of an additional application on another one of the clusters with the additional application utilizing remote data resources locally accessible within the other cluster.

11. The method of claim 10 wherein the application master component is configured to access a resolving interface of the distributed catalog service in order to determine for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource or a remote data resource.

12. The method of claim 10 wherein the application master component comprises:

a local cluster node manager;

one or more remote cluster node managers; and

an application programming interface through which at least one of a client and the aggregator communicates with the local and remote cluster node managers.

13. The method of claim 12 wherein the local cluster node manager of the application master component is configured to control operation of the aggregator of the application master component.

14. A method comprising:

initiating a first application in a first one of a plurality of distributed processing node clusters associated with respective data zones, each of the clusters being configured to perform processing operations utilizing local data resources locally accessible within its corresponding data zone;

determining a plurality of data resources to be utilized by the application;

identifying for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource that is locally accessible within the data zone of the first distributed processing node cluster or a remote data resource that is not locally accessible within the data zone of the first distributed processing node cluster;

for one or more of the plurality of data resources that are identified as local data resources, performing processing operations utilizing the local data resources in the first cluster in accordance with the first application;

for one or more of the plurality of data resources that are identified as remote data resources, initiating respective additional applications in one or more additional ones of the plurality of distributed processing node clusters and performing processing operations utilizing the remote data resources in the corresponding one or more additional clusters in accordance with the one or more additional applications;

aggregating processing results from the first and one or more additional clusters; and

providing the aggregated processing results to a client;

wherein the method further comprises accessing a distributed catalog service to identify for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource or a remote data resource;

wherein the distributed catalog service is distributed over the clusters with each of the clusters having visibility of a corresponding distinct portion of the distributed catalog based on its locally accessible data resources; and

wherein the method is implemented by at least one processing device comprising a processor coupled to a memory.

15. 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 at least one processing device causes said at least one processing device:

to initiate a first application in a first one of a plurality of distributed processing node clusters associated with respective data zones, each of the clusters being configured to perform processing operations utilizing local data resources locally accessible within its corresponding data zone;

to determine a plurality of data resources to be utilized by the application;

to identify for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource that is locally accessible within the data zone of the first distributed processing node cluster or a remote data resource that is not locally accessible within the data zone of the first distributed processing node cluster;

for one or more of the plurality of data resources that are identified as local data resources, to perform processing operations utilizing the local data resources in the first cluster in accordance with the first application;

for one or more of the plurality of data resources that are identified as remote data resources, to initiate respective additional applications in one or more additional ones of the plurality of distributed processing node clusters and to perform processing operations utilizing the remote data resources in the corresponding one or more additional clusters in accordance with the one or more additional applications;

to aggregate processing results from the first and one or more additional clusters; and

to provide the aggregated processing results to a client;

wherein the processing results from the first and one or more additional clusters preserve at least one specified policy of those clusters in the respective local and remote data resources wherein the specified policy relates to at least one of privacy, security, governance, risk and compliance; and

wherein the processing results from a given one of the clusters are permitted to be transmitted to another one of the clusters but the local data resources of the given cluster that are utilized to obtain the processing results are not permitted to be transmitted to another one of the clusters.

16. The computer program product of claim 15 wherein a given one of the applications initiated in a corresponding one of the clusters comprises:

an application master component;

at least one aggregator controlled by the application master component; and

one or more containers utilized by the aggregator in conjunction with execution of the given application;

wherein the one or more containers correspond to respective ones of the local data resources that are locally accessible within the corresponding cluster; and

wherein the aggregator is configured to request initiation of an additional application on another one of the clusters with the additional application utilizing remote data resources locally accessible within the other cluster.

17. The computer program product of claim 15 wherein the program code when executed by said at least one processing device further causes said at least one processing device to access a distributed catalog service to identify for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource or a remote data resource, wherein the distributed catalog service is distributed over the clusters with each of the clusters having visibility of a corresponding distinct portion of the distributed catalog based on its locally accessible data resources.

18. An apparatus comprising:

at least one processing device having a processor coupled to a memory;

wherein said at least one processing device is configured:

to initiate a first application in a first one of a plurality of distributed processing node clusters associated with respective data zones, each of the clusters being configured to perform processing operations utilizing local data resources locally accessible within its corresponding data zone;

to determine a plurality of data resources to be utilized by the application;

to identify for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource that is locally accessible within the data zone of the first distributed processing node cluster or a remote data resource that is not locally accessible within the data zone of the first distributed processing node cluster;

for one or more of the plurality of data resources that are identified as local data resources, to perform processing operations utilizing the local data resources in the first cluster in accordance with the first application;

for one or more of the plurality of data resources that are identified as remote data resources, to initiate respective additional applications in one or more additional ones of the plurality of distributed processing node clusters and to perform processing operations utilizing the remote data resources in the corresponding one or more additional clusters in accordance with the one or more additional applications;

to aggregate processing results from the first and one or more additional clusters; and

to provide the aggregated processing results to a client;

wherein the processing results from the first and one or more additional clusters preserve at least one specified policy of those clusters in the respective local and remote data resources wherein the specified policy relates to at least one of privacy, security, governance, risk and compliance; and

wherein the processing results from a given one of the clusters are permitted to be transmitted to another one of the clusters but the local data resources of the given cluster that are utilized to obtain the processing results are not permitted to be transmitted to another one of the clusters.

19. An information processing system comprising the plurality of distributed processing node clusters of claim 18 .

20. The apparatus of claim 18 wherein said at least one processing device is further configured to access a distributed catalog service to identify for each of the plurality of data resources to be utilized by the application whether the data resource is a local data resource or a remote data resource, wherein the distributed catalog service is distributed over the clusters with each of the clusters having visibility of a corresponding distinct portion of the distributed catalog based on its locally accessible data resources.

Assignments (10)
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 (045455/0001) Recorded May 20, 2022
From: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS NOTES COLLATERAL AGENT
To: DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO ASAP SOFTWARE EXPRESS, INC.); DELL MARKETING L.P. (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO CREDANT TECHNOLOGIES, INC.); DELL USA L.P.; DELL INTERNATIONAL L.L.C.; DELL PRODUCTS L.P.; DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO FORCE10 NETWORKS, INC. AND WYSE TECHNOLOGY L.L.C.); EMC CORPORATION (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO MAGINATICS LLC); EMC IP HOLDING COMPANY LLC (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO MOZY, INC.); SCALEIO LLC
Reel/Frame 061753/0001 →
RELEASE OF SECURITY INTEREST IN PATENTS PREVIOUSLY RECORDED AT REEL/FRAME (040136/0001) Recorded Apr 26, 2022
From: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS NOTES COLLATERAL AGENT
To: DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO ASAP SOFTWARE EXPRESS, INC.); DELL MARKETING L.P. (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO CREDANT TECHNOLOGIES, INC.); DELL USA L.P.; DELL INTERNATIONAL L.L.C.; DELL PRODUCTS L.P.; DELL MARKETING CORPORATION (SUCCESSOR-IN-INTEREST TO FORCE10 NETWORKS, INC. AND WYSE TECHNOLOGY L.L.C.); EMC CORPORATION (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO MAGINATICS LLC); EMC IP HOLDING COMPANY LLC (ON BEHALF OF ITSELF AND AS SUCCESSOR-IN-INTEREST TO MOZY, INC.); SCALEIO LLC
Reel/Frame 061324/0001 →
RELEASE OF SECURITY INTEREST Recorded Nov 3, 2021
From: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH
To: ASAP SOFTWARE EXPRESS, INC.; AVENTAIL LLC; CREDANT TECHNOLOGIES, INC.; DELL USA L.P.; DELL INTERNATIONAL, L.L.C.; DELL MARKETING L.P.; DELL PRODUCTS L.P.; DELL SOFTWARE INC.; DELL SYSTEMS CORPORATION; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; FORCE10 NETWORKS, INC.; MAGINATICS LLC; MOZY, INC.; SCALEIO LLC; WYSE TECHNOLOGY L.L.C.
Reel/Frame 058216/0001 →
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 Sep 29, 2016
From: EMC CORPORATION
To: EMC IP HOLDING COMPANY LLC
Reel/Frame 040203/0001 →
SECURITY AGREEMENT Recorded Sep 21, 2016
From: ASAP SOFTWARE EXPRESS, INC.; AVENTAIL LLC; CREDANT TECHNOLOGIES, INC.; DELL USA L.P.; DELL INTERNATIONAL L.L.C.; DELL MARKETING L.P.; DELL PRODUCTS L.P.; DELL SOFTWARE INC.; DELL SYSTEMS CORPORATION; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; FORCE10 NETWORKS, INC.; MAGINATICS LLC; MOZY, INC.; SCALEIO LLC; SPANNING CLOUD APPS LLC; WYSE TECHNOLOGY L.L.C.
To: THE BANK OF NEW YORK MELLON TRUST COMPANY, N.A., AS NOTES COLLATERAL AGENT
Reel/Frame 040136/0001 →
SECURITY AGREEMENT Recorded Sep 21, 2016
From: ASAP SOFTWARE EXPRESS, INC.; AVENTAIL LLC; CREDANT TECHNOLOGIES, INC.; DELL USA L.P.; DELL INTERNATIONAL L.L.C.; DELL MARKETING L.P.; DELL PRODUCTS L.P.; DELL SOFTWARE INC.; DELL SYSTEMS CORPORATION; EMC CORPORATION; EMC IP HOLDING COMPANY LLC; FORCE10 NETWORKS, INC.; MAGINATICS LLC; MOZY, INC.; SCALEIO LLC; SPANNING CLOUD APPS LLC; WYSE TECHNOLOGY L.L.C.
To: CREDIT SUISSE AG, CAYMAN ISLANDS BRANCH, AS COLLATERAL AGENT
Reel/Frame 040134/0001 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 12, 2016
From: FLORISSI, PATRICIA GOMES SOARES; LUTATI, BENNY; GUDES, EHUD; GONEN, YARON; SINGER, IDO; MEISELS, AMNON; VIJENDRA, SUDHIR
To: EMC CORPORATION
Reel/Frame 038258/0204 →
Continuity (2)
Provisional Application 62143404 · Apr 6, 2015
Provisional Application 62143685 · Apr 6, 2015
Cited By (7)
US 12,235,736 US 12,309,070 US 12,339,750 US 12,489,657 US 12,505,002 US 12,561,125 US 12,675,368