IP Library Granted Patent US 12,222,915
Granted Patent B2
US 12,222,915 · App. 17/314,813 · Granted Feb 11, 2025

Accumulating and flushing mutations in a column store

Inventor: Todd Lipcon (San Francisco, CA)
Assignee: CLOUDERA, INC.
G06F16/221G06F16/2246G06F16/2282G06F16/23
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 12,222,915
App. No.
17/314,813
Granted
Feb 11, 2025
Kind
B2
Abstract

Columnar storage provides many performance and space saving benefits for analytic workloads, but previous mechanisms for handling single row update transactions in column stores suffer from poor performance. A columnar data layout facilitates both low-latency random access capabilities together with high-throughput analytical access capabilities, simplifying Hadoop architectures for use cases involving real-time data. In disclosed embodiments, mutations within a single row are executed atomically across columns and do not necessarily include the entirety of a row. This allows for faster updates without the overhead of reading or rewriting larger columns.

Claims (53)

1. A method for a distributed database system, the method comprising:

in response to an insert operation for a database table, storing a new row in a first in-memory store, wherein the first in-memory store is configured to flush to an on-disk store after a flushing criteria is reached;

prior to the first in-memory store flushing the new row to the on-disk store, storing a separate mutation record linked to the new row that is stored in the first in-memory store, in response to a second mutation operation that is specified for the new row; and

after the flushing criteria is reached, performing:

modifying the new row stored in the first in-memory store based on applying the separate mutation record to the new row; and

flushing (i) the modified new row to one or more base data files in the on-disk store, and (ii) the separate mutation record associated with the new row to one or more delta files in the on-disk store,

wherein the first in-memory store and the on-disk store are configured for disjoint row storage such that a given row of the database table exists in one of a row set (RowSet) in the first in-memory store or a RowSet in the on-disk store.

2. The method of claim 1 , wherein the separate mutation record is flushed to the one or more delta files with a timestamp that is associated with the second mutation operation.

3. The method of claim 1 , further comprising:

in response to a third mutation operation that specifies a flushed row stored in the on-disk store, storing a mutation record in a second in-memory store that is associated with the on-disk store, wherein the second in-memory store is configured to flush to the one or more delta files in the on-disk store after a second flushing criteria is reached.

4. The method of claim 1 ,

wherein the one or more base data files and the one or more delta files are immutable.

5. The method of claim 1 ,

wherein the database table includes (1) a number of columns, including a primary key (PK) column, and (2) a number of rows, each row having a unique PK stored in its corresponding PK column, wherein the number of rows are divided into a plurality of tablets, each tablet divided into a number of RowSets,

wherein the RowSets are disjoint with respect to a stored key so that any given key is active in at most one RowSet, and

wherein, for a given PK, at most one RowSet contains an active copy of a row with the given PK.

6. The method of claim 1 , wherein the database table remains readable by a client of the database table during the flushing of the modified new row and the flushing of the separate mutation record.

7. The method of claim 1 , further comprising:

subsequent to flushing the modified new row, storing a new set of rows of the database table in the first in-memory store.

8. The method of claim 1 , further comprising:

processing a time-travel query for the modified new row in the on-disk store based on accessing the separate mutation record in the one or more delta files in the on-disk store.

9. A computing system comprising:

one or more processors; and

one or more memories storing instructions that, when executed by the one or more processors, cause the computing system to perform a process comprising:

in response to a first insert operation for a database table, storing a new row in a first in-memory store, wherein the first in-memory store is configured to flush to an on-disk store after a flushing criteria is reached;

prior to the first in-memory store flushing the new row to the on-disk store, storing a separate mutation record attached to the new row that is stored in the first in-memory store, in response to a second mutation operation that is specified for the new row; and

after the flushing criteria is reached, performing:

modifying the new row stored in the first in-memory store based on applying the separate mutation record to the new row; and

flushing (i) the modified new row to one or more base data files in the on-disk store, and (ii) the separate mutation record associated with the new row to one or more delta files in the on-disk store,

wherein the first in-memory store and the on-disk store are configured for disjoint row storage such that a given row of the database table exists in one of a row set (RowSet) in the first in-memory store or a RowSet in the on-disk store.

10. The computing system of claim 9 , wherein the process further comprises:

determining that a third mutation operation specifies a flushed row that is stored in the on-disk store based on a Bloom filter for a set of primary keys present in the on-disk store; and

storing a mutation record for the third mutation operation in a second in-memory store instead of the first in-memory store based on the third mutation operation specifying the flushed row stored in the on-disk store, wherein the second in-memory store is configured to flush to the one or more delta files in the on-disk store after a second flushing criteria is reached.

11. The computing system of claim 10 , wherein the process further comprises:

chunking available Bloom filters into pages corresponding to ranges of rows;

creating an index for the pages using an immutable B-tree; and

storing the pages and the index in a server-wide least recent used (LRU) page cache.

12. The computing system of claim 9 ,

wherein the separate mutation record attached to the new row is implemented as a singly linked list,

wherein the singly linked list includes a mutation head that points to a given mutation node, and a next mutation pointer in the given mutation node points to a subsequent mutation node.

13. A non-transitory computer-readable medium storing instructions that, when executed by a computing system, cause the computing system to perform a process comprising:

in response to a first insert operation for a database table, storing a new row in a first in-memory store, wherein the first in-memory store is configured to flush to an on-disk store after a flushing criteria is reached;

prior to the first in-memory store flushing the new row to the on-disk store, storing a separate mutation record associated with the new row that is stored in the first in-memory store, in response to a second mutation operation that is specified for the new row; and

after the flushing criteria, performing:

modifying the new row stored in the first in-memory store based on applying the separate mutation record to the new row; and

flushing (i) the modified new row to one or more base data files in the on-disk store, and (ii) flushing the separate mutation record associated with the new row to one or more delta files in the on-disk store,

wherein the first in-memory store and the on-disk store are configured for disjoint row storage such that a given row of the database table exists in one of a row set (RowSet) in the first in-memory store or a RowSet in the on-disk store.

14. The non-transitory computer-readable medium of claim 13 , wherein the process further comprises:

performing a minor delta compaction operation that reduces a number of one or more delta files in the on-disk store by merging the one or more delta files together, without updating the one or more base data files in the on-disk store.

15. The non-transitory computer-readable medium of claim 13 , wherein the separate mutation record is a REDO record, and wherein the process further comprises:

performing a major delta compaction operation that migrates REDO records to UNDO records and updatess the one or more base data files in the on-disk store.

16. The non-transitory computer-readable medium of claim 13 , wherein the process further comprises:

performing a merging compaction operation that merges together a number of RowSets, in one or more on-disk stores, that overlap in range.

Assignments (3)
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 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded May 26, 2021
From: LIPCON, TODD
To: CLOUDERA, INC
Reel/Frame 056393/0793 →
Continuity (4)
Continuation 15881541 · Jan 26, 2018
Continuation 15149128 · May 7, 2016
Provisional Application 62158444 · May 7, 2015
Related Publication 20210271653A1 · Sep 2, 2021
References Cited (93)
US 6026406A · Huang · 2000 [cited by examiner]
US 6119128A · Courter · 2000 [cited by examiner]
US 6349310B1 · Klein et al. · 2002 [cited by applicant]
US 7548928B1 · Dean · 2009 [cited by applicant]
US 7567973B1 · Burrows et al. · 2009 [cited by applicant]
US 8010492B2 · Biswal et al. · 2011 [cited by applicant]
US 8452737B2 · Netz et al. · 2013 [cited by applicant]
US 8706769B1 · Gao · 2014 [cited by examiner]
US 9280570B2 · Pruner · 2016 [cited by applicant]
US 9519676B1 · Finnie · 2016 [cited by examiner]
US 9747169B2 · Kottomtharayil et al. · 2017 [cited by applicant]
US 10346432B2 · Lipcon · 2019 [cited by applicant]
US 11003642B2 · Lipcon · 2021 [cited by applicant]
US 20030018644A1 · Bala et al. · 2003 [cited by applicant]
US 20030065884A1 · Lu · 2003 [cited by examiner]
US 20040148303A1 · Mckay et al. · 2004 [cited by applicant]
US 20040199519A1 · Gu · 2004 [cited by examiner]
US 20050240615A1 · Barsness et al. · 2005 [cited by applicant]
US 20070033354A1 · Burrows et al. · 2007 [cited by applicant]
US 20110010330A1 · McCline · 2011 [cited by examiner]
US 20110022574A1 · Hansen · 2011 [cited by applicant]
US 20110082854A1 · Eidson · 2011 [cited by applicant]
US 20110213775A1 · Franke · 2011 [cited by applicant]
US 20110258242A1 · Eidson · 2011 [cited by examiner]
US 20110276744A1 · Sengupta et al. · 2011 [cited by applicant]
US 20120016851A1 · Hrle et al. · 2012 [cited by applicant]
US 20120016852A1 · Hrle et al. · 2012 [cited by applicant]
US 20120072656A1 · Archak et al. · 2012 [cited by applicant]
US 20120078978A1 · Shoolman et al. · 2012 [cited by applicant]
US 20120221528A1 · Renkes et al. · 2012 [cited by applicant]
US 20120323971A1 · Pasupuleti · 2012 [cited by applicant]
US 20130117247A1 · Schreter · 2013 [cited by examiner]
US 20130166554A1 · Yoon et al. · 2013 [cited by applicant]
US 20130198139A1 · Skidanov · 2013 [cited by examiner]
US 20130218840A1 · Smith · 2013 [cited by examiner]
US 20130275476A1 · Chandler · 2013 [cited by applicant]
US 20130290282A1 · Faerber · 2013 [cited by examiner]
US 20130311612A1 · Dickinson · 2013 [cited by applicant]
US 20130339293A1 · Witten et al. · 2013 [cited by applicant]
US 20140081918A1 · Srivas · 2014 [cited by examiner]
US 20140136788A1 · Faerber · 2014 [cited by examiner]
US 20140164823A1 · Zheng · 2014 [cited by examiner]
US 20140195492A1 · Wilding · 2014 [cited by examiner]
US 20140215170A1 · Scarpino et al. · 2014 [cited by applicant]
US 20140244599A1 · Zhang et al. · 2014 [cited by applicant]
US 20140279855A1 · Tan et al. · 2014 [cited by applicant]
US 20140297601A1 · Pruner · 2014 [cited by applicant]
US 20140304219A1 · Yoon · 2014 [cited by applicant]
US 20140317048A1 · Wang · 2014 [cited by examiner]
US 20150026128A1 · Drobychev et al. · 2015 [cited by applicant]
US 20150039573A1 · Bhattacharjee · 2015 [cited by examiner]
US 20150046413A1 · Andrei et al. · 2015 [cited by applicant]
US 20150088824A1 · Kamp · 2015 [cited by examiner]
US 20150213071A1 · Alvey · 2015 [cited by examiner]
US 20150237127A1 · Khemani · 2015 [cited by examiner]
US 20150242400A1 · Bensberg · 2015 [cited by examiner]
US 20150278281A1 · Zhang · 2015 [cited by applicant]
US 20150286668A1 · Legler · 2015 [cited by applicant]
US 20160042016A1 · Faerber et al. · 2016 [cited by applicant]
US 20160070618A1 · Pundir et al. · 2016 [cited by applicant]
US 20160275094A1 · Lipcon · 2016 [cited by examiner]
US 20160328429A1 · Lipcon et al. · 2016 [cited by applicant]
US 20170351718A1 · Faerber · 2017 [cited by examiner]
US 20180150490A1 · Lipcon · 2018 [cited by applicant]
US 20180253468A1 · Gurajada · 2018 [cited by examiner]
CN 1737943A · 2006 [cited by examiner]
CN 104111962B · 2018 [cited by examiner]
EP 2898432B1 · 2020 [cited by applicant]
JP 2010146113A · 2010 [cited by applicant]
WO 2013074665 · 2013 [cited by applicant]
H. Li, “Flash Saver: Save the Flash-Based Solid State Drives through Deduplication and Delta-encoding,” 2012 13th International Conference on Parallel and Distributed Computing, Applications and Technologies, Beijing, C… [cited by examiner]
Heman et al., “Read-Optimized Data Storage”, Jul. 2008. (Year: 2008). [cited by examiner]
Otlu, Süleyman Onur. A New Technique : Replace Algorithm to Retrieve a Version from a Repository Instead of Delta Application. Middle East Technical University, 2004. (Year: 2004). [cited by examiner]
Data Engineering, IEEE Computer Society, vol. 36 No. 2, Jun. 2013. (Year: 2013). [cited by examiner]
Jagatheesan, A., Levandoski, J., Neumann, T., Pavlo, A. (eds) In Memory Data Management and Analysis. IMDM IMDM 2013 2014 . Lecture Notes in Computer Science(), vol. 8921. Springer, Cham. https://doi.org/10.1007/978-3-3… [cited by examiner]
U.S. Appl. No. 15/073,509, Lipcon, et al. [cited by applicant]
Yang , et al., “A Distributed Storage Model for EHR Based on HBase”, 2011 International Conference on Information ManagementInnovation Management and Industrial Engineering. [cited by applicant]
Chang, F. et al., “Bigtable: A Distributed Storage System for Structured Data.” Google, Inc., OSDI 2006, 14 pages. [cited by applicant]
Extended Search Report for European Application No. 16790218.8; Date of Mailing: Jan. 4, 2018, 6 pages. [cited by applicant]
Extended Search Report for European Application No. 16765788.1; Date of Mailing: Jan. 12, 2018, 9 pages. [cited by applicant]
Examination Report for European Application No. 16765788.1; Date of Mailing: Apr. 16, 2019; 6 pages. [cited by applicant]
Lipcon et al., “Kudu: Storage for Fast Analytics on Fast Data*,” Sep. 28, 2015; retrieved from the internet: URL:https://kudu.apache.org/kudo.pdf [retrieved on Dec. 13, 2017]; 13 pages. [cited by applicant]
International Search Report and Written Opinion for Application No. PCT/US2016/031391; Applicant: Cloudera, Inc.; Date of Mailing Aug. 17, 2016, 11 pages. [cited by applicant]
International Search Report and Written Opinion for International Application No. PCT/US2016/022968; Applicant: Cloudera, Inc.; Date of Mailing: Jul. 8, 2016, 7 pages. [cited by applicant]
Menon, et al., “Optimizing Key-Value Stores for Hybrid Storage Architectures,” CASCON '14 Proceedings of the 24th Annual International Conference on Computer Science and Software Engineering, Nov. 3, 2014; pp. 355-358. [cited by applicant]
Sakr et al., A Survey of Large Scale Data Management Approaches in Cloud Environments, IEEE Communications Surveys & Tutorials, vol. 13, No. 3, Third Quarter 2011, pp. 311-336. [cited by applicant]
Zeng, J. et al., “Multi-Tenant Fair Share in NoSQL Data Stores,” IEEE International Conference on Cluster Computing (CLUSTER), Sep. 22-26, 2014, pp. 176-184. [cited by applicant]
Hofhansl, L., About HBase flushes and compactions, hadoop-hbase.blogspot.com/2014/07/about-hbase-flushes-and-compactions.html, pp. 1-2. (Year: 2014). [cited by applicant]
Wei et al., i BigTable: Practical Data Integrity for BigTable in Public Cloud, CODASPY'13, Feb. 18-20, 2013, San Antonio, Texas, USA, pp. 341-352. (Year: 2013). [cited by applicant]
Wikipedia, “Multiversion concurrency control”, May 1, 2015 version (Year: 2015), 6 pages. [cited by applicant]
Bernstein et al., “Concurrency Control in Distributed Database Systems”, Jun. 1981 (Year: 1981), 37 pages. [cited by applicant]
Wei , et al., “iBigTable: Practical Data Integrity for BigTable in Public Cloud”, CODASPY 2013, USA, Feb. 18-20, 2013, pp. 341-352. [cited by applicant]
Burns et al., In-place reconstruction of delta compressed files, 1998, Proceedings of the seventeenth annual ACM symposium on principles of distributed computing (PODC'98, Association for computing machinery, New York, … [cited by applicant]