IP Library › Granted Patent US 12,210,899
Granted Patent B1
US 12,210,899 · App. 17/346,887 · Granted Jan 28, 2025

Processing data payloads from external data storage systems in an observability pipeline system

Inventors: Dritan Bitincka (Edgewater, NJ); Ledion Bitincka (San Francisco, CA); Nicholas Robert Romito (Chicago, IL); Clint Sharp (Oakland, CA)
Assignee: Cribl, Inc.
G06F9/4881G06F9/3851G06F9/44505G06F9/541
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,210,899
App. No.
17/346,887
Granted
Jan 28, 2025
Kind
B1
Abstract

Data payloads from an external data storage system are processed in an observability pipeline system. In some aspects, the observability pipeline system defines a leader role and worker roles. The leader role generates a data discovery task based on configuration information for a data collection task. A worker role executes the data discovery task, which includes communicating with an external data storage system to identify a data payload that is stored on the external data storage system and contains event data that meet event filter criteria. The leader role generates data collection tasks based on the data payload. Worker roles execute the data collection tasks. Executing a data collection task includes: communicating with the external data storage system to obtain a subset of filtered event data from the data payload; and streaming the subset of filtered event data to an observability pipeline process.

Claims (70)

1. A method of operating an observability pipeline system, the method comprising:

in an observability pipeline system deployed on one or more computer nodes operating as a leader role and a plurality of worker roles, receiving configuration information for a data collection job, the configuration information comprising event filter criteria;

by operation of the leader role, generating a data discovery task based on the configuration information;

by operation of one of the worker roles, executing the data discovery task, wherein executing the data discovery task comprises communicating with an external data storage system to identify a data payload that is stored on the external data storage system and contains event data that meet the event filter criteria;

by operation of the leader role, generating a plurality of data collection tasks based on the data payload identified by the execution of the data discovery task;

by operation of one or more of the worker roles, executing the plurality of data collection tasks, wherein executing each respective data collection task comprises:

communicating with the external data storage system to obtain a subset of filtered event data from the data payload, each subset of filtered event data comprising a respective portion of the event data that meet the event filter criteria; and

streaming the subset of filtered event data to an observability pipeline process.

2. The method of claim 1 , wherein the configuration information is received from a user device, and the event filter criteria are based on input received through a user interface of the user device.

3. The method of claim 2 , wherein the observability pipeline process is selected through the user interface of the user device and indicated in the configuration information from the user device.

4. The method of claim 1 , wherein the data payload comprises a set of files, each of the plurality of data collection tasks identifies a respective one of the files, and executing a data collection task comprises communicating with the external data storage system to obtain a subset of filtered event data from the file identified by the data collection task.

5. The method of claim 1 , wherein the data payload comprises a set of files, a first data collection task of the plurality of data collection tasks identifies multiple files, and executing the first data collection task comprises communicating with the external data storage system to obtain a subset of filtered event data from the multiple files identified by the first data collection task.

6. The method of claim 1 , comprising applying observability pipeline processes to the respective subsets of filtered event data, wherein applying an observability pipeline process to a subset of filtered event data comprises:

applying schema normalization to the subset of filtered event data to generate normalized event data;

routing the normalized event data to a streaming analytics and processing engine;

generating structured data from the normalized event data by operation of the streaming analytics and processing engine; and

applying one or more output schemas to the structured data to generate observability pipeline output data for one or more external data destinations.

7. The method of claim 1 , wherein the observability pipeline system is deployed on a distributed computer system comprising a leader node operating as the leader role and a plurality of worker nodes operating as the worker roles.

8. The method of claim 1 , wherein the observability pipeline system is deployed on a standalone computer system comprising a single computer node operating as the leader role and the plurality of worker roles.

9. The method of claim 1 , wherein the data collection job is a full run mode data collection job, and the method comprises, prior to receiving the configuration information for the full run mode data collection job:

in the observability pipeline system, receiving configuration information for a preview mode data collection job;

by operation of the leader role, generating a preview data discovery task based on the configuration information for the preview mode data collection job;

by operation of one of the worker roles, executing the preview data discovery task;

by operation of the leader role, generating a plurality of preview data collection tasks based on a data payload identified by the execution of the preview data discovery task;

by operation of one or more of the worker roles, executing the plurality of preview data collection tasks, wherein executing each respective preview data collection task comprises communicating with the external data storage system to obtain a subset of filtered event data from the data payload identified by the execution of the preview data discovery task; and

sending, from the observability pipeline system to a user device, the filtered event data obtained by the execution of the plurality of preview data collection tasks.

10. The method of claim 1 , comprising, prior to receiving the configuration information for the data collection job:

at the observability pipeline system, receiving pipeline input data comprising event data from a plurality of data sources;

by operation of the observability pipeline system, generating pipeline output data by applying one or more observability pipeline processes to the event data from the plurality of data sources; and

delivering the pipeline output data to a plurality of external data destinations, wherein delivering the pipeline output data comprises storing the data payload on the external data storage system.

11. An observability pipeline system comprising one or more computer nodes that operate a leader role and a plurality of worker roles, the one or more computer nodes comprising:

one or more processors; and

memory storing instructions that, when executed by the one or more processors, cause the one or more processors to perform operations comprising:

receiving configuration information for a data collection job, the configuration information comprising event filter criteria;

by operation of the leader role, generating a data discovery task based on the configuration information;

by operation of one of the worker roles, executing the data discovery task, wherein executing the data discovery task comprises communicating with an external data storage system to identify a data payload that is stored on the external data storage system and contains event data that meet the event filter criteria;

by operation of the leader role, generating a plurality of data collection tasks based on the data payload identified by the execution of the data discovery task;

by operation of one or more of the worker roles, executing the plurality of data collection tasks, wherein executing each respective data collection task comprises:

communicating with the external data storage system to obtain a subset of filtered event data from the data payload, each subset of filtered event data comprising a respective portion of the event data that meet the event filter criteria; and

streaming the subset of filtered event data to an observability pipeline process.

12. The system of claim 11 , comprising a communication interface configured to receive the configuration information from a user device, and the event filter criteria are based on input received through a user interface of the user device.

13. The system of claim 11 , wherein the data payload comprises a set of files, each of the plurality of data collection tasks identifies a respective one of the files, and executing a data collection task comprises communicating with the external data storage system to obtain a subset of filtered event data from the file identified by the data collection task.

14. The system of claim 11 , wherein the data payload comprises a set of files, a first data collection task of the plurality of data collection tasks identifies multiple files, and executing the first data collection task comprises communicating with the external data storage system to obtain a subset of filtered event data from the multiple files identified by the first data collection task.

15. The system of claim 11 , the operations comprising applying observability pipeline processes to the respective subsets of filtered event data, wherein applying an observability pipeline process to a subset of filtered event data comprises:

applying schema normalization to the subset of filtered event data to generate normalized event data;

routing the normalized event data to a streaming analytics and processing engine;

generating structured data from the normalized event data by operation of the streaming analytics and processing engine; and

applying one or more output schemas to the structured data to generate observability pipeline output data for one or more external data destinations.

16. The system of claim 11 , wherein the observability pipeline system comprises a distributed computer system comprising a leader node operating as the leader role and a plurality of worker nodes operating as the worker roles.

17. The system of claim 11 , wherein the observability pipeline system comprises a standalone computer system comprising a single computer node operating as the leader role and the plurality of worker roles.

18. The system of claim 11 , wherein the data collection job is a full run mode data collection job, and the operations comprise, prior to receiving the configuration information for the full run mode data collection job:

receiving configuration information for a preview mode data collection job;

by operation of the leader role, generating a preview data discovery task based on the configuration information for the preview mode data collection job;

by operation of one of the worker roles, executing the preview data discovery task;

by operation of the leader role, generating a plurality of preview data collection tasks based on a data payload identified by the execution of the preview data discovery task;

by operation of one or more of the worker roles, executing the plurality of preview data collection tasks, wherein executing each respective preview data collection task comprises communicating with the external data storage system to obtain a subset of filtered event data from the data payload identified by the execution of the preview data discovery task; and

sending, from the observability pipeline system to a user device, the filtered event data obtained by the execution of the plurality of preview data collection tasks.

19. The system of claim 11 , the operations comprising, prior to receiving the configuration information for the data collection job:

receiving pipeline input data comprising event data from a plurality of data sources;

by operation of the observability pipeline system, generating pipeline output data by applying one or more observability pipeline processes to the event data from the plurality of data sources; and

delivering the pipeline output data to a plurality of external data destinations, wherein delivering the pipeline output data comprises storing the data payload on the external data storage system.

20. A non-transitory computer-readable medium comprising instructions that are operable when executed by data processing apparatus to perform operations comprising:

defining a leader role and a plurality of worker roles of an observability pipeline system;

receiving configuration information for a data collection job, the configuration information comprising event filter criteria;

by operation of the leader role, generating a data discovery task based on the configuration information;

by operation of one of the worker roles, executing the data discovery task, wherein executing the data discovery task comprises communicating with an external data storage system to identify a data payload that is stored on the external data storage system and contains event data that meet the event filter criteria;

by operation of the leader role, generating a plurality of data collection tasks based on the data payload identified by the execution of the data discovery task;

by operation of one or more of the worker roles, executing the plurality of data collection tasks, wherein executing each respective data collection task comprises:

communicating with the external data storage system to obtain a subset of filtered event data from the data payload, each subset of filtered event data comprising a respective portion of the event data that meet the event filter criteria; and

streaming the subset of filtered event data to an observability pipeline process.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jun 14, 2021
From: BITINCKA, DRITAN; BITINCKA, LEDION; ROMITO, NICHOLAS ROBERT; SHARP, CLINT
To: CRIBL, INC.
Reel/Frame 056536/0001 →
References Cited (27)
US 10108605B1 · Leighton · 2018 [cited by examiner]
US 10367676B1 · Vermeulen · 2019 [cited by examiner]
US 10691728B1 · Masson et al. · 2020 [cited by applicant]
US 11057414B1 · Giorgio · 2021 [cited by examiner]
US 12067419B1 · Bitincka et al. · 2024 [cited by applicant]
US 20120066683A1 · Srinath · 2012 [cited by applicant]
US 20180300174A1 · Karanasos et al. · 2018 [cited by applicant]
US 20180349212A1 · Liu et al. · 2018 [cited by applicant]
US 20190121978A1 · Kraemer · 2019 [cited by examiner]
US 20200272114A1 · Grabowski · 2020 [cited by examiner]
US 20200327037A1 · Toal et al. · 2020 [cited by applicant]
US 20210182869A1 · Davis · 2021 [cited by examiner]
Scrocca et al. The Kaiju Project: Enabling Event-Driven Observability. [online] (Jul. 17). ACM., pp. 85-96. Retrieved From the Internet <https://dl.acm.org/doi/pdf/10.1145/3401025.3401740> (Year: 2020). [cited by examiner]
Sambasivan et al. Principled workflow-centric tracing of distributed systems. [online] (Oct. 7). ACM., pp. 401-414. Retrieved From the Internet <https://dl.acm.org/doi/pdf/10.1145/2987550.2987568> (Year: 2016). [cited by examiner]
“Load Balancing (Computing)”, https://en.wikipedia.org/w/index.php?title=Load_balancing_(computing) &oldid=1028284441, Jun. 13, 2021, 15 pgs. [cited by applicant]
“Round-robin Scheduling”, https://en.wikipedia.org/w/index.php?title=Round-robin_scheduling&oldid=1019633430, Apr. 24, 2021, 5 pgs. [cited by applicant]
Bitincka, Dritan , “Collectors”, https://docs.cribl.io/docs/collectors, accessed Jun. 13, 2021, version last updated Jun. 5, 2021, 6 pgs. [cited by applicant]
Cribl, Inc. , “Distributed Deployment”, https://docs.cribl.io/docs/deploy-distributed, accessed Jun. 13, 2021, last updated May 21, 2021, 33 pgs. [cited by applicant]
Litras, Steve , “Data Collection is Here”, https://cribl.io/blog/data-collection-is-here/ accessed Jun. 13, 2021, dated Jun. 15, 2020, 7 pgs. [cited by applicant]
Litras , “Working with Data in LogStream 2.2”, https://cribl.io/blog/working-with-data-in-logstream-2-2/, Jul. 14, 2020, 5 pgs. [cited by applicant]
Romito , “Demystifying Collection Job Scheduling”, https://cribl.io/blog/demystifying-collection-job-scheduling/, Jun. 24, 2020, 9 pgs. [cited by applicant]
Sharp , “Building an observability pipeline on top of open source Apache NiFi, Logstash, or Fluentd: a journey”, https://cribl.io/blog/building-an-observability-pipeline-on-top-of-open-source-apache-nifi-logstash-or-flu… [cited by applicant]
Sharp, Clint , “Grappling with Observability Data Management”, https://thenewstack.io/grappling-with-observability-data-management/, Mar. 31, 2021, 13 pgs. [cited by applicant]
Sharp , “The Observability Pipeline”, https://cribl.io/blog/the-observability-pipeline/, Oct. 10, 2019, 10 pgs. [cited by applicant]
Treat, Tyler , “Microservice Observability, Part 2: Evolutionary Patterns for Solving Observability Problems”, https://bravenewgeek.com/microservice-observability-part-2-evolutionary-patterns-for-solving-observability-p… [cited by applicant]
Turiff, Bryan , “Announcing Cribl LogStream 2.2: Baby Got Batch!”, https://cribl.io/logstream-2-2-baby-got-batch/ accessed Jun. 13, 2021, dated Jun. 15, 2020, 9 pgs. [cited by applicant]
USPTO, Notice of Allowance issued in U.S. Appl. No. 17/346,881 on Jun. 1, 2023, 14 pages. [cited by applicant]
Cited By (2)
US 12,470,372 US 12,476,793