IP Library Granted Patent US 11,567,956
Granted Patent B2
US 11,567,956 · App. 16/921,640 · Granted Jan 31, 2023

Background format optimization for enhanced queries in a distributed computing cluster

Inventors: Marcel Kornacker (Oakland, CA); Justin Erickson (San Francisco, CA); Nong Li (San Francisco, CA); Lenni Kuff (San Francisco, CA); Henry Noel Robinson (San Francisco, CA); Alan Choi (Palo Alto, CA); Alex Behm (San Francisco, CA)
Assignee: Cloudera, Inc.
G06F16/2471G06F16/24542G06F16/258G06F16/27
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,567,956
App. No.
16/921,640
Granted
Jan 31, 2023
Kind
B2
Abstract

A format conversion engine for Apache Hadoop that converts data from its original format to a database-like format at certain time points for use by a low latency (LL) query engine. The format conversion engine comprises a daemon that is installed on each data node in a Hadoop cluster. The daemon comprises a scheduler and a converter. The scheduler determines when to perform the format conversion and notifies the converter when the time comes. The converter converts data on the data node from its original format to a database-like format for use by the low latency (LL) query engine.

Claims (38)

1. A method for performing queries on stored data in a distributed computing cluster, the method comprising:

creating, by a query engine at a first data node in a plurality of data nodes forming a peer-to-peer network, a query fragment based on an availability of converted data at the first data node that interacts with other data nodes, the converted data corresponding to data associated with a query received from a client, wherein the converted data is the data associated with the query converted, by the first data node, from an original format into a target format that is specified by a schema;

causing, by the query engine at the first data node, execution of the query fragment on data that corresponds to a format for which the query fragment is created, based on the schema;

obtaining, by the query engine at the first data node, an intermediate result based on the execution of the query fragment; and

aggregating, by the query engine at the first data node for the client, the intermediate result with other intermediate results executed at query engines at the other data nodes in the peer-to-peer network.

2. The method of claim 1 , wherein the target format is a columnar format.

3. The method of claim 1 , wherein the target format is configured for relational database processing.

4. The method of claim 1 , wherein, when the converted data is available, the query fragment is created for the target format.

5. The method of claim 1 , wherein, when the converted data is not available, the query fragment is created for the original format.

6. The method of claim 1 , further comprising:

parsing the query from the client;

wherein the query fragment is created in response to parsing the query from the client.

7. The method of claim 1 , further comprising:

forming, based on the plurality of data nodes, the peer-to-peer network;

wherein each data node functions as a peer in the peer-to-peer network.

8. The method of claim 1 , further comprising:

distributing, by the first data node, the query fragment to a second data node of the plurality of data nodes.

9. The method of claim 1 , wherein the distributed computing cluster is a Hadoop™ cluster.

10. A computer system for performing queries on stored data in a distributed computing cluster, the computer system comprising:

one or more processors; and

one or more memory modules having instructions stored thereon, which when executed by the one or more processors, cause the computer system to:

create, by a query engine at a first data node in a plurality of data nodes forming a peer-to-peer network, a query fragment based on an availability of converted data at the first data node that interacts with other data nodes, the converted data corresponding to data associated with a query received from a client, wherein the converted data is the data associated with the query converted, by the first data node, from an original format into a target format that is specified by a schema;

cause, by the query engine at the first data node, execution of the query fragment on data that corresponds to a format for which the query fragment is created, based on the schema;

obtain, by the query engine at the first data node, an intermediate result based on the execution of the query fragment; and

aggregate, by the query engine at the first data node for the client, the intermediate result with other intermediate results from the other data nodes in the distributed computing cluster.

11. The computer system of claim 10 , wherein the target format is a columnar format.

12. The computer system of claim 10 , wherein the target format is configured for relational database processing.

13. The computer system of claim 10 , wherein, when the converted data is available, the query fragment is created for the target format.

14. The computer system of claim 10 , wherein, when the converted data is not available, the query fragment is created for the original format.

15. A non-transitory machine-readable storage medium having stored thereon instructions which, when executed by one or more processors, cause the one or more processors to perform a method comprising:

creating, by a query engine at a first data node in a plurality of data nodes forming a peer-to-peer network, a query fragment based on an availability of converted data at the first data node that interacts with other data nodes, the converted data corresponding to data associated with a query received from a client, wherein the converted data is the data associated with the query converted, by the first data node, from an original format into a target format that is specified by a schema;

causing, by the query engine at the first data node, execution of the query fragment on data that corresponds to a format for which the query fragment is created, based on the schema;

obtaining, by the query engine at the first data node, an intermediate result based on the execution of the query fragment; and

aggregating, by the query engine at the first data node for the client, the intermediate result with other intermediate results from the other data nodes in a distributed computing cluster.

16. The non-transitory machine-readable storage medium of claim 15 , wherein the target format is a columnar format.

17. The non-transitory machine-readable storage medium of claim 15 , wherein the target format is configured for relational database processing.

18. The non-transitory machine-readable storage medium of claim 15 , wherein, when the converted data is available, the query fragment is created for the target format.

19. The non-transitory machine-readable storage medium of claim 15 , wherein, when the converted data is not available, the query fragment is created for the original format.

Assignments (5)
RELEASE OF SECURITY INTERESTS IN PATENTS Recorded Oct 14, 2021
From: CITIBANK, N.A.
To: CLOUDERA, INC.; HORTONWORKS, INC.
Reel/Frame 057804/0355 →
FIRST LIEN NOTICE AND CONFIRMATION OF GRANT OF SECURITY INTEREST IN PATENTS Recorded Oct 12, 2021
From: CLOUDERA, INC.; HORTONWORKS, INC.
To: JPMORGAN CHASE BANK, N.A.
Reel/Frame 057776/0185 →
SECOND LIEN NOTICE AND CONFIRMATION OF GRANT OF SECURITY INTEREST IN PATENTS Recorded Oct 12, 2021
From: CLOUDERA, INC.; HORTONWORKS, INC.
To: JPMORGAN CHASE BANK, N.A.
Reel/Frame 057776/0284 →
SECURITY INTEREST Recorded Dec 22, 2020
From: CLOUDERA, INC.; HORTONWORKS, INC.
To: CITIBANK, N.A., AS COLLATERAL AGENT
Reel/Frame 054832/0559 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jul 7, 2020
From: KORNACKER, MARCEL; ERICKSON, JUSTIN; LI, NONG; KUFF, LENNI; ROBINSON, HENRY NOEL; CHOI, ALAN; BEHM, ALEX
To: CLOUDERA, INC.
Reel/Frame 053137/0916 →