IP Library Granted Patent US 10,706,059
Granted Patent B2
US 10,706,059 · App. 15/292,053 · Granted Jul 7, 2020

Background format optimization for enhanced SQL-like queries in Hadoop

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 10,706,059
App. No.
15/292,053
Granted
Jul 7, 2020
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 system for performing queries on stored data in a Hadoop™ distributed computing cluster, the system comprising:

a plurality of data nodes forming a peer-to-peer network for the queries received from a client, a respective data node of the plurality of data nodes functioning as a peer in the peer-to-peer network and being capable of interacting with components of the Hadoop™ cluster, the respective data node operating an instance of a query engine that is configured to:

parse a query from a client;

selectively creates query fragments based on an availability of converted data at the respective data node, the converted data corresponding to data associated with the query, wherein the converted data is the data associated with the query converted from an original format into a target format that is specified by a schema, and wherein the query is processed into said query fragments by whichever data node that receives the query;

distribute the query fragments among the plurality of data nodes;

execute the query fragments on whichever local data that corresponds to a format for which the query fragments are created, based on the schema;

obtain intermediate results from other data nodes that receive the query fragments; and

aggregate the intermediate results for the client.

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

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

4. The system of claim 1 , wherein, when the converted data is available, the query fragments are created for the target format.

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

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

forming, based on a plurality of data nodes, a peer-to-peer network for the queries received from a client, a respective data node of the plurality of data nodes functioning as a peer in the peer-to-peer network and being capable of interacting with components of the Hadoop™ cluster; and

in the respective data node, operating an instance of a query engine that is configured to:

parse a query from a client;

selectively creates query fragments based on an availability of converted data at the respective data node, the converted data corresponding to data associated with the query, wherein the converted data is the data associated with the query converted from an original format into a target format that is specified by a schema, and wherein the query is processed into said query fragments by whichever data node that receives the query;

distribute the query fragments among the plurality of data nodes;

execute the query fragments on whichever local data that corresponds to a format for which the query fragments are created, based on the schema;

obtain intermediate results from other data nodes that receive the query fragments; and

aggregate the intermediate results for the client.

7. The method of claim 6 , wherein the target format is a columnar format.

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

9. The method of claim 6 , wherein, when the converted data is available, the query fragments are created for the target format.

10. The method of claim 6 , wherein, when the converted data is not available, the query fragments are created for the original format.

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

forming, based on a plurality of data nodes, a peer-to-peer network for the queries received from a client, a respective data node of the plurality of data nodes functioning as a peer in the peer-to-peer network and being capable of interacting with components of a Hadoop™ cluster; and

in the respective data node, operating an instance of a query engine that is configured to:

parse a query from a client;

selectively creates query fragments based on an availability of converted data at the respective data node, the converted data corresponding to data associated with the query, wherein the converted data is the data associated with the query converted from an original format into a target format that is specified by a schema, and wherein the query is processed into said query fragments by whichever data node that receives the query;

distribute the query fragments among the plurality of data nodes;

execute the query fragments on whichever local data that corresponds to a format for which the query fragments are created, based on the schema;

obtain intermediate results from other data nodes that receive the query fragments; and

aggregate the intermediate results for the client.

12. The medium of claim 11 , wherein the target format is a columnar format.

13. The medium of claim 11 , wherein the target format is configured for relational database processing.

14. The medium of claim 11 , wherein, when the converted data is available, the query fragments are created for the target format.

15. The medium of claim 11 , wherein, when the converted data is not available, the query fragments are 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 Oct 14, 2016
From: KORNACKER, MARCEL; ERICKSON, JUSTIN; LI, NONG; KUFF, LENNI; ROBINSON, HENRY NOEL; CHOI, ALAN; BEHM, ALEX
To: CLOUDERA, INC.
Reel/Frame 040361/0616 →
Continuity (2)
Continuation 14043753 · Oct 1, 2013
Related Publication 20170032003A1 · Feb 2, 2017