IP Library Granted Patent US 12,339,859
Granted Patent B2
US 12,339,859 · App. 17/664,774 · Granted Jun 24, 2025

User interface for managing distributed query execution

Inventors: David C. Tracey (Monasterboice, IE); Miguel A. Casanova (Dublin, IE)
Assignee: Rapid7, Inc.
G06F16/2471G06F11/3006G06F16/1734G06F16/182G06F16/2322G06F16/242
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,339,859
App. No.
17/664,774
Granted
Jun 24, 2025
Kind
B2
Abstract

Systems and methods are disclosed to implement a distributed query execution system that performs statistical operations on specified time windows over time-based datasets. In embodiments, the query system splits a statistical function into a set of parallel accumulator tasks that correspond to different portions of the dataset and/or function time windows. The accumulator tasks are executed in parallel by individual accumulator nodes to generate individual statistical result structures. The structures are then combined by an aggregator node to produce an aggregate result structure that indicates the results of the statistical function over the time windows. In embodiments, the accumulator and aggregator tasks are implemented and executed using a programmable task execution framework that allows developers to define custom accumulator and aggregator tasks. Advantageously, the query system allows queries with time-windowed statistical functions to be parallelized across a group of worker nodes and scaled to very large datasets.

Claims (51)

1. A method, comprising:

performing, by one or more hardware processors with associated memory that implement a distributed query execution system:

storing a time-based dataset on one or more storage devices, wherein the time-based dataset is stored as a plurality of files on a plurality of servers;

receiving, via a user interface, a query directed to the time-based dataset, wherein the query specifies (a) a user-specified statistical function to be computed over groups of records in the time-based dataset and (b) a user-specified number and size of a plurality of time windows over which to compute the statistical function;

executing the query using a plurality of compute nodes connected via a network, including:

executing a set of accumulator nodes in parallel to read respective portions of the time-based dataset in individual ones of the files and compute the statistic function over the respective portion, wherein individual ones of the accumulator nodes are launched on respective servers that have local access to different ones of the files so that each respective portion is locally read and processed by a respective group of one or more accumulator nodes at a respective server so as to reduce data traffic over the network;

executing at least one aggregator node to aggregate result structures produced by the accumulator nodes;

detecting that one of the accumulator nodes has failed prior to completion of computation; and

retrying the computation of the failed accumulator node using a different accumulator node;

responsive to user input received via the user interface, pausing the execution of the query by the compute nodes;

while the execution of the query is paused, dynamically updating the user interface to output partial results of the statistical function based on a subset of the accumulator nodes that have completed computation;

responsive to additional user input received via the user interface after the pausing of the query, resuming the execution of the query; and

dynamically updating the user interface to output complete results of the statistical function after completion of the execution of the query.

2. The method of claim 1 , wherein the user interface is a graphical user interface.

3. The method of claim 2 , wherein the partial results of the statistical function are displayed on the graphical user interface as a time graph.

4. The method of claim 1 , wherein the compute nodes are virtual machine instances hosted on one or more virtual machine hosts.

5. The method of claim 1 , wherein the accumulator nodes include two or more accumulator nodes that share a common cache or file lock when reading the files.

6. The method of claim 1 , wherein the distributed query execution system is configured to compute a plurality of statistical functions on time-based datasets, including two or more of:

a count of matched records,

a byte size of matched records,

a total value of an attribute in matched records,

a minimum of an attribute in matched records,

a maximum of an attribute in matched records, and

an average of an attribute in matched records.

7. A system comprising:

a distributed query execution system implemented by one or more hardware processors with associated memory, configure to:

store a time-based dataset on one or more storage devices, wherein the time-based dataset is stored as a plurality of files on a plurality of servers;

receive, via a user interface, a query directed to the time-based dataset, wherein the query specifies (a) a user-specified statistical function to be computed over groups of records in the time-based dataset and (b) a user-specified number and size of a plurality of time windows over which to compute the statistical function;

execute the query using a plurality of compute nodes connected via a network, including to:

execute a set of accumulator nodes in parallel to read respective portions of the time-based dataset in individual ones of the files and compute the statistic function over the respective portion, wherein individual ones of the accumulator nodes are launched on respective servers that have local access to different ones of the files so that each respective portion is locally read and processed by a respective group of one or more accumulator nodes at a respective server so as to reduce data traffic over the network;

execute at least one aggregator node to aggregate result structures produced by the accumulator nodes;

detect that one of the accumulator nodes has failed prior to completion of computation; and

retry the computation of the failed accumulator node using a different accumulator node;

responsive to user input received via the user interface, pause the execution of the query by the compute nodes;

while the execution of the query is paused, dynamically update the user interface to output partial results of the statistical function based on a subset of the accumulator nodes that have completed computation;

responsive to additional user input received via the user interface after the pause of the query, resume the execution of the query; and

dynamically update the user interface to output complete results of the statistical function after completion of the execution of the query.

8. The system of claim 7 , wherein the user interface is a graphical user interface.

9. The system of claim 8 , wherein the partial results of the statistical function are displayed on the graphical user interface as a time graph.

10. One or more non-transitory computer-accessible storage media storing program instructions that when executed on or across one or more processors implement at least a portion of a distributed query execution system and cause the distributed query execution system to:

store a time-based dataset on one or more storage devices, wherein the time-based dataset is stored as a plurality of files on a plurality of servers;

receive, via a user interface, a query directed to the time-based dataset, wherein the query specifies (a) a user-specified statistical function to be computed over groups of records in the time-based dataset and (b) a user-specified number and size of a plurality of time windows over which to compute the statistical function;

execute the query using a plurality of compute nodes connected via a network, including to:

execute a set of accumulator nodes in parallel to read respective portions of the time-based dataset in individual ones of the files and compute the statistic function over the respective portion, wherein individual ones of the accumulator nodes are launched on respective servers that have local access to different ones of the files so that each respective portion is locally read and processed by a respective group of one or more accumulator nodes at a respective server so as to reduce data traffic over the network;

execute at least one aggregator node to aggregate result structures produced by the accumulator nodes;

detect that one of the accumulator nodes has failed prior to completion of computation; and

retry the computation of the failed accumulator node using a different accumulator node;

responsive to user input received via the user interface, pause the execution of the query by the compute nodes;

while the execution of the query is paused, dynamically update the user interface to output partial results of the statistical function based on a subset of the accumulator nodes that have completed computation;

responsive to additional user input received via the user interface after the pause of the query, resume the execution of the query; and

dynamically update the user interface to output complete results of the statistical function after completion of the execution of the query.

Assignments (2)
SECURITY INTEREST Recorded Jun 26, 2025
From: RAPID7, INC.; RAPID7 LLC
To: JPMORGAN CHASE BANK, N.A.
Reel/Frame 071743/0537 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 6, 2022
From: TRACEY, DAVID C; CASANOVA, MIGUEL A
To: RAPID7, INC.
Reel/Frame 061332/0798 →
Continuity (2)
Continuation 16798222 · Feb 21, 2020
Related Publication 20220284019A1 · Sep 8, 2022
References Cited (82)
US 6470335B1 · Marusak · 2002 [cited by examiner]
US 6959320B2 · Shah · 2005 [cited by examiner]
US 7620633B1 · Parsons et al. · 2009 [cited by applicant]
US 7779340B2 · Milne et al. · 2010 [cited by applicant]
US 7814482B2 · Kwon et al. · 2010 [cited by applicant]
US 8090730B2 · Shahabi et al. · 2012 [cited by applicant]
US 8260803B2 · Hsu et al. · 2012 [cited by applicant]
US 8510807B1 · Elazary et al. · 2013 [cited by applicant]
US 9251035B1 · Vazac et al. · 2016 [cited by applicant]
US 9277003B2 · Doyle et al. · 2016 [cited by applicant]
US 9607067B2 · Haas · 2017 [cited by examiner]
US 9646024B2 · Srivas et al. · 2017 [cited by applicant]
US 9923918B2 · Nicodemus · 2018 [cited by examiner]
US 10152558B2 · Lisobee et al. · 2018 [cited by applicant]
US 10176246B2 · Dang · 2019 [cited by examiner]
US 10248621B2 · Ward et al. · 2019 [cited by applicant]
US 10296500B2 · Dean et al. · 2019 [cited by applicant]
US 10296658B2 · Le Biannic et al. · 2019 [cited by applicant]
US 10579634B2 · Erdogan et al. · 2020 [cited by applicant]
US 10747736B2 · Klemenz · 2020 [cited by examiner]
US 10776355B1 · Batsakis · 2020 [cited by examiner]
US 10855712B2 · Oliner et al. · 2020 [cited by applicant]
US 11093526B2 · Shmueli · 2021 [cited by examiner]
US 11182434B2 · Beedgen · 2021 [cited by examiner]
US 11372871B1 · Tracey · 2022 [cited by examiner]
US 20020065628A1 · McMahan · 2002 [cited by examiner]
US 20030115183A1 · Abdo · 2003 [cited by examiner]
US 20040066741A1 · Dinker · 2004 [cited by examiner]
US 20050065921A1 · Hrle · 2005 [cited by examiner]
US 20070271547A1 · Gulko et al. · 2007 [cited by applicant]
US 20080059489A1 · Han · 2008 [cited by examiner]
US 20090228446A1 · Anzai · 2009 [cited by examiner]
US 20090254971A1 · Herz · 2009 [cited by examiner]
US 20100057684A1 · Williamson · 2010 [cited by examiner]
US 20120011134A1 · Travnik · 2012 [cited by examiner]
US 20120054172A1 · Agrawal · 2012 [cited by examiner]
US 20120078951A1 · Hsu et al. · 2012 [cited by applicant]
US 20120159258A1 · Maybee · 2012 [cited by examiner]
US 20120191699A1 · George et al. · 2012 [cited by applicant]
US 20130007003A1 · Shyr et al. · 2013 [cited by applicant]
US 20130198165A1 · Cheng et al. · 2013 [cited by applicant]
US 20140040276A1 · Chen · 2014 [cited by examiner]
US 20140172867A1 · Lin · 2014 [cited by examiner]
US 20140337274A1 · Unnikrishnan · 2014 [cited by examiner]
US 20140351233A1 · Crupi et al. · 2014 [cited by applicant]
US 20150149441A1 · Nica · 2015 [cited by examiner]
US 20160285805A1 · Oliver et al. · 2016 [cited by applicant]
US 20160328432A1 · Raghunathan · 2016 [cited by applicant]
US 20160335318A1 · Gerweck et al. · 2016 [cited by applicant]
US 20160371363A1 · Muro et al. · 2016 [cited by applicant]
US 20170220938A1 · Sainani et al. · 2017 [cited by applicant]
US 20180052897A1 · Namarvar · 2018 [cited by examiner]
US 20180089258A1 · Bhattacharjee et al. · 2018 [cited by applicant]
US 20180246950A1 · Arye · 2018 [cited by examiner]
US 20190050453A1 · Duffield · 2019 [cited by examiner]
US 20200050612A1 · Bhattacharjee · 2020 [cited by examiner]
US 20200177562A1 · Tav · 2020 [cited by examiner]
US 20200186455A1 · Lokhandwala · 2020 [cited by examiner]
US 20200387818A1 · Chan et al. · 2020 [cited by applicant]
US 20210250306A1 · Gladney · 2021 [cited by examiner]
EP 0325081A2 · 1989 [cited by examiner]
EP 2876562A1 · 2015 [cited by examiner]
WO 2008144732 · 2008 [cited by applicant]
WO WO2011003232A1 · 2011 [cited by examiner]
WO 2013144535 · 2013 [cited by applicant]
WO 2014039336 · 2014 [cited by applicant]
WO 2014144889 · 2014 [cited by applicant]
WO WO2014169265A1 · 2014 [cited by examiner]
WO 2016165509 · 2016 [cited by applicant]
WO WO2018170276A2 · 2018 [cited by examiner]
WO WO2020033446A1 · 2020 [cited by examiner]
WO WO2020227645A1 · 2020 [cited by examiner]
Wenhuchen et al., “ADatasetforAnswering Time-SensitiveQuestions”, 35thConferenceonNeuralInformationProcessingSystems(NeurIPS2021)TrackonDatasetsandBenchmarks, Oct. 2021, pp. 1-17. [cited by examiner]
Umar Albalawi, “Countermeasure of Statistical Inference in Database Security”, 2018 IEEE International Conference on Big Data (Big Data), Dec. 2018, pp. 2044-2047. [cited by examiner]
H. Andrade et al., “Optimizing the execution of multiple data analysis queries on parallel and distributed environments”, IEEE Transactions on Parallel and Distributed Systems (vol. 15, Issue: 6, 2004, pp. 520-532). [cited by examiner]
Yuan Wei et al., “Maintaining data freshness in distributed real-time databases”, Proceedings. 16th Euromicro Conference on Real-Time Systems, 2004. ECRTS 2004. (2004, pp. 251-260). [cited by examiner]
Saket Sathe et al., “AFFINITY: Efficiently querying statistical measures on time-series data”, IEEE 29th International Conference on Data Engineering (ICDE), Apr. 2013, pp. 841-852. [cited by applicant]
Liang Zhang et al., “TARDIS: Distributed Indexing Framework for Big Time Series Data”, IEEE 35th International Conference on Data Engineering (ICDE), Jun. 2019, pp. 1202-1213. [cited by applicant]
Pietro Pinoli et al., “Metadata management for scientific databases”, Information Systems 81, 2019, pp. 1-20. [cited by applicant]
Natasha Noy et al., “Google Dataset Search Building a search engine for datasets in an open Web ecosystem”, WWW 19: The World Wide Web Conference, May 2019, pp. 1365-1375. [cited by applicant]
H. Papageorgiou et al., “Modeling statistical metadata”, Proceedings Thirteenth International Conference on Scientific and Statistical Database Management, SSDBM 2001, Jul. 2001, pp. 25-35. [cited by applicant]
U.S. Appl. No. 16/798,222, filed Feb. 21, 2020, David C. Tracey. [cited by applicant]