MLchartDataset catalogue

Patent · US10019308B1 · B1 · US

Disaster-proof event data processing

(11) Publication number
US10019308B1
(21) Application number
14/860,441
(22) Filing date
2015-09-21
(30) Priority date
2012-01-31
(43) Publication date
2018-07-10
(45) Date of grant
2018-07-10
(51) IPC
G06F 11/07
(52) CPC
  • G06F Electric digital data processing: 11/0793, 11/0709, 11/0751, 11/182, 11/187, 11/202
  • H04L Transmission of digital information, e.g. telegraphic communication: 51/02, 51/23
(73) Assignee
Google LLC
(72) Inventors
Ashish Gupta; Haifeng Jiang; Manpreet Singh; Monica Chawathe
(54) Title
Disaster-proof event data processing
(57) Abstract

Systems and methods are disclosed herein for providing fault tolerant processing of events. The system includes multiple consensus computers configured to communicate with one another and multiple event processors configured to process data such as events. Each consensus computer is further configured to receive a request to process a unit of data from an event processor. A consensus computer communicates with at least one other consensus computer to reach consensus as to whether the unit of data has previously been assigned to an event processor for processing. Then, a consensus computer sends a message to the event processor that sent the inquiry including instructions to either process the unit of data or not process the unit of data. Because the consensus computers determine whether a unit of data has previously been assigned to an event processor, the system ensures that an event is not processed more than once.

Full text
View on Google Patents

Claims (20)

  1. A method of providing fault tolerant event processing, comprising: receiving, by a particular consensus computer from among a plurality of consensus computers in communication with one another, a request to process a unit of data from a particular event processor from among a plurality of event processors, wherein each of the plurality of consensus computers stores metadata in its own metadata database regarding which units of data have been processed by that consensus computer; forwarding, by the particular consensus computer, the request to process the unit of data from the particular event processor to the plurality of consensus computers; determining, by the plurality of consensus computers, whether the unit of data is to be processed by the particular event processor based on a determination of a consensus of whether any event processor of the plurality of event processors has already committed to processing the unit of data based on the metadata stored in each of the plurality of consensus computer's respective metadata database; and instructing, by the particular consensus computer, the particular event processor on how to process the unit of data in accordance with the consensus.
  2. The method of claim 1, wherein the consensus is determined based on applying a Paxos consensus algorithm to one or more votes received from the plurality of consensus computers.
  3. The method of claim 1, wherein the unit of data includes bundles of one or more action events.
  4. The method of claim 3, further comprising: storing information indicative of the one or more action events on a consensus system.
  5. The method of claim 1, wherein each consensus computer is configured to: vote on whether a plurality of consensus computers previously assigned each action event included in a bundle to be processed by an event processor; and instruct the event processor that sent the request to process the unit of data to process a portion of the bundle or not process the bundle based on the consensus.
  6. The method of claim 1, wherein each consensus computer is configured to: count votes from at least one other consensus computer of the plurality of consensus computers; and determine whether the number of votes indicating that the plurality of consensus computers have previously assigned the unit of data to be processed by an event processor exceeds a threshold, wherein the instructions are based on the determination.
  7. The method of claim 1, wherein the plurality of consensus computers determines the consensus using a fault tolerant consensus algorithm that tolerates the failure of a consensus computer of the plurality of consensus computers or the failure of an event processor.
  8. The method of claim 1, wherein the plurality of consensus computers are configured to: determine whether any event processor of the plurality of event processors is currently processing the unit of data; and responsive to determining that an event processor of the plurality of event processors is currently processing the unit of data, instructing the particular event processor not to process the unit of data.
  9. The method of claim 8, wherein the particular event processor and the particular consensus computer are located at separate geographic locations.
  10. The method of claim 1, further comprising: based on a determination the particular event processor is approved to process the unit of data by the plurality of consensus computers, updating, by the particular consensus computer, the particular consensus computer's metadata database to reflect that the event processor has to committed to process the unit of data.
  11. A consensus system comprising: a plurality of consensus computers in communication with each other, each consensus computer comprising storage units storing instructions that when executed by a particular consensus computer of the plurality of consensus computers cause the particular consensus computer to perform operations comprising: receiving a request to process a unit of data from a particular event processor of a plurality of event processors, wherein each of the plurality of consensus computers stores metadata in its own metadata database regarding which units of data have been processed by that consensus computer; forwarding, the request to process the unit of data from the particular event processor to the plurality of consensus computers; determining whether the unit of data is to be processed by the particular event processor based on a determination from the plurality of consensus computers of a consensus of whether any event processor of the plurality of event processors has committed to processing the unit of data based on the metadata stored in each of the plurality of consensus computer's respective metadata database; and instructing the particular event processor on how to process the unit of data in accordance with the consensus.
  12. The system of claim 11, wherein the consensus is determined based on applying a Paxos consensus algorithm to one or more votes received from the plurality of consensus computers.
  13. The system of claim 11, wherein the unit of data includes bundles of one or more action events.
  14. The system of claim 13, the operations further comprising: storing information indicative of the one or more action events.
  15. The system of claim 11, wherein each consensus computer in the plurality of consensus computers is configured to: vote on whether one or more consensus computers previously assigned each action event included in a bundle to be processed by an event processor; and instruct the particular event processor that sent the request to process the unit of data to process a portion of the bundle or not process the bundle based on the consensus.
  16. The system of claim 11, wherein each consensus computer in the plurality of consensus computers is configured to: count votes from at least one other consensus computer in the plurality of consensus computers; and determine whether the number of votes indicating that one or more consensus computers have previously assigned the unit of data to be processed by an event processor exceeds a threshold, wherein the instructions are based on the determination.
  17. The system of claim 11, wherein the plurality of consensus computers determines the consensus using a fault tolerant consensus algorithm that tolerates the failure of a consensus computer of the plurality of consensus computers or the failure of an event processor.
  18. The system of claim 11, wherein each consensus computer in the plurality of consensus computers is configured to: determine whether any event processor of a plurality of event processors is currently processing the unit of data; and responsive to determining that an event processor of the plurality of event processors is currently processing the unit of data, instructing the particular event processor not to process the unit of data.
  19. The system of claim 18, wherein the particular event processor and the particular consensus computer are located at separate geographic locations.
  20. The system of claim 11, the operations further comprising: based on a determination the event processor is approved to process the unit of data by the plurality of consensus computers, updating the particular consensus computer's metadata database to reflect that the event processor has to committed to process the unit of data.

Description

In general, the systems and methods disclosed herein describe reliable and fault tolerant processing of event data.

Hosts of websites sometimes sell or lease space on their websites to content providers. Often, it is useful for the hosts to provide content relevant to the content providers, such as data related to traffic on the website. In particular, the content providers often pay for space on a website based on the number of interactions (mouse clicks or mouseovers) users have with such websites. The data may be related to user clicks, queries, views, or any other type of user interaction with the website. Sometimes, systems that process this data fail or need to be temporarily disabled for maintenance, and data is lost. This can lead to inaccuracies in the reported statistics and can also lead content providers to underpay or overpay the website host. Thus, robust mechanisms are desired for providing reliable aggregate data in a useful form.

Accordingly, systems and methods disclosed herein provide reliable, fault tolerant processing of online user events. According to one aspect, the disclosure relates to a system for processing event data. The system comprises multiple consensus computers configured to communicate with one another. Each consensus computer is further configured to receive an inquiry from an event processor. The inquiry includes a request to process a unit of data. A consensus computer communicates with at least one other consensus computer to reach consensus as to whether the consensus computers previously assigned the unit of data to be processed by an event processor.

Citations (18)

  • US6463532B1
  • US20010025351A1
  • US20100017495A1
  • US20050149609A1
  • US20050256824A1
  • US20050283659A1
  • US20050283373A1
  • US20060069942A1
  • US20060136781A1
  • US20060168011A1
  • US7376867B1
  • US20060156312A1
  • US20070214355A1
  • US20090260012A1
  • US20100017644A1
  • US20100082728A1
  • US20140189270A1
  • US9172670B1
Record as JSON
{
  "publication_number": "US10019308B1",
  "country": "US",
  "kind": "B1",
  "title": "Disaster-proof event data processing",
  "abstract": "Systems and methods are disclosed herein for providing fault tolerant processing of events. The system includes multiple consensus computers configured to communicate with one another and multiple event processors configured to process data such as events. Each consensus computer is further configured to receive a request to process a unit of data from an event processor. A consensus computer communicates with at least one other consensus computer to reach consensus as to whether the unit of data has previously been assigned to an event processor for processing. Then, a consensus computer sends a message to the event processor that sent the inquiry including instructions to either process the unit of data or not process the unit of data. Because the consensus computers determine whether a unit of data has previously been assigned to an event processor, the system ensures that an event is not processed more than once.",
  "claims": [
    "1. A method of providing fault tolerant event processing, comprising: receiving, by a particular consensus computer from among a plurality of consensus computers in communication with one another, a request to process a unit of data from a particular event processor from among a plurality of event processors, wherein each of the plurality of consensus computers stores metadata in its own metadata database regarding which units of data have been processed by that consensus computer; forwarding, by the particular consensus computer, the request to process the unit of data from the particular event processor to the plurality of consensus computers; determining, by the plurality of consensus computers, whether the unit of data is to be processed by the particular event processor based on a determination of a consensus of whether any event processor of the plurality of event processors has already committed to processing the unit of data based on the metadata stored in each of the plurality of consensus computer's respective metadata database; and instructing, by the particular consensus computer, the particular event processor on how to process the unit of data in accordance with the consensus.",
    "2. The method of claim 1, wherein the consensus is determined based on applying a Paxos consensus algorithm to one or more votes received from the plurality of consensus computers.",
    "3. The method of claim 1, wherein the unit of data includes bundles of one or more action events.",
    "4. The method of claim 3, further comprising: storing information indicative of the one or more action events on a consensus system.",
    "5. The method of claim 1, wherein each consensus computer is configured to: vote on whether a plurality of consensus computers previously assigned each action event included in a bundle to be processed by an event processor; and instruct the event processor that sent the request to process the unit of data to process a portion of the bundle or not process the bundle based on the consensus.",
    "6. The method of claim 1, wherein each consensus computer is configured to: count votes from at least one other consensus computer of the plurality of consensus computers; and determine whether the number of votes indicating that the plurality of consensus computers have previously assigned the unit of data to be processed by an event processor exceeds a threshold, wherein the instructions are based on the determination.",
    "7. The method of claim 1, wherein the plurality of consensus computers determines the consensus using a fault tolerant consensus algorithm that tolerates the failure of a consensus computer of the plurality of consensus computers or the failure of an event processor.",
    "8. The method of claim 1, wherein the plurality of consensus computers are configured to: determine whether any event processor of the plurality of event processors is currently processing the unit of data; and responsive to determining that an event processor of the plurality of event processors is currently processing the unit of data, instructing the particular event processor not to process the unit of data.",
    "9. The method of claim 8, wherein the particular event processor and the particular consensus computer are located at separate geographic locations.",
    "10. The method of claim 1, further comprising: based on a determination the particular event processor is approved to process the unit of data by the plurality of consensus computers, updating, by the particular consensus computer, the particular consensus computer's metadata database to reflect that the event processor has to committed to process the unit of data.",
    "11. A consensus system comprising: a plurality of consensus computers in communication with each other, each consensus computer comprising storage units storing instructions that when executed by a particular consensus computer of the plurality of consensus computers cause the particular consensus computer to perform operations comprising: receiving a request to process a unit of data from a particular event processor of a plurality of event processors, wherein each of the plurality of consensus computers stores metadata in its own metadata database regarding which units of data have been processed by that consensus computer; forwarding, the request to process the unit of data from the particular event processor to the plurality of consensus computers; determining whether the unit of data is to be processed by the particular event processor based on a determination from the plurality of consensus computers of a consensus of whether any event processor of the plurality of event processors has committed to processing the unit of data based on the metadata stored in each of the plurality of consensus computer's respective metadata database; and instructing the particular event processor on how to process the unit of data in accordance with the consensus.",
    "12. The system of claim 11, wherein the consensus is determined based on applying a Paxos consensus algorithm to one or more votes received from the plurality of consensus computers.",
    "13. The system of claim 11, wherein the unit of data includes bundles of one or more action events.",
    "14. The system of claim 13, the operations further comprising: storing information indicative of the one or more action events.",
    "15. The system of claim 11, wherein each consensus computer in the plurality of consensus computers is configured to: vote on whether one or more consensus computers previously assigned each action event included in a bundle to be processed by an event processor; and instruct the particular event processor that sent the request to process the unit of data to process a portion of the bundle or not process the bundle based on the consensus.",
    "16. The system of claim 11, wherein each consensus computer in the plurality of consensus computers is configured to: count votes from at least one other consensus computer in the plurality of consensus computers; and determine whether the number of votes indicating that one or more consensus computers have previously assigned the unit of data to be processed by an event processor exceeds a threshold, wherein the instructions are based on the determination.",
    "17. The system of claim 11, wherein the plurality of consensus computers determines the consensus using a fault tolerant consensus algorithm that tolerates the failure of a consensus computer of the plurality of consensus computers or the failure of an event processor.",
    "18. The system of claim 11, wherein each consensus computer in the plurality of consensus computers is configured to: determine whether any event processor of a plurality of event processors is currently processing the unit of data; and responsive to determining that an event processor of the plurality of event processors is currently processing the unit of data, instructing the particular event processor not to process the unit of data.",
    "19. The system of claim 18, wherein the particular event processor and the particular consensus computer are located at separate geographic locations.",
    "20. The system of claim 11, the operations further comprising: based on a determination the event processor is approved to process the unit of data by the plurality of consensus computers, updating the particular consensus computer's metadata database to reflect that the event processor has to committed to process the unit of data."
  ],
  "description_excerpt": "In general, the systems and methods disclosed herein describe reliable and fault tolerant processing of event data.\n\nHosts of websites sometimes sell or lease space on their websites to content providers. Often, it is useful for the hosts to provide content relevant to the content providers, such as data related to traffic on the website. In particular, the content providers often pay for space on a website based on the number of interactions (mouse clicks or mouseovers) users have with such websites. The data may be related to user clicks, queries, views, or any other type of user interaction with the website. Sometimes, systems that process this data fail or need to be temporarily disabled for maintenance, and data is lost. This can lead to inaccuracies in the reported statistics and can also lead content providers to underpay or overpay the website host. Thus, robust mechanisms are desired for providing reliable aggregate data in a useful form.\n\nAccordingly, systems and methods disclosed herein provide reliable, fault tolerant processing of online user events. According to one aspect, the disclosure relates to a system for processing event data. The system comprises multiple consensus computers configured to communicate with one another. Each consensus computer is further configured to receive an inquiry from an event processor. The inquiry includes a request to process a unit of data. A consensus computer communicates with at least one other consensus computer to reach consensus as to whether the consensus computers previously assigned the unit of data to be processed by an event processor.",
  "cpc": [
    "G06F 11/0793",
    "G06F 11/0709",
    "G06F 11/0751",
    "G06F 11/182",
    "G06F 11/187",
    "G06F 11/202",
    "H04L 51/02",
    "H04L 51/23"
  ],
  "ipc": [
    "G06F 11/07"
  ],
  "assignees": [
    "Google LLC"
  ],
  "inventors": [
    "Ashish Gupta",
    "Haifeng Jiang",
    "Manpreet Singh",
    "Monica Chawathe"
  ],
  "filing_date": "2015-09-21",
  "publication_date": "2018-07-10",
  "grant_date": "2018-07-10",
  "priority_date": "2012-01-31",
  "application_number": "US-201514860441-A",
  "family_id": "54328340",
  "cited_by_count": 2,
  "citations": [
    "US6463532B1",
    "US20010025351A1",
    "US20100017495A1",
    "US20050149609A1",
    "US20050256824A1",
    "US20050283659A1",
    "US20050283373A1",
    "US20060069942A1",
    "US20060136781A1",
    "US20060168011A1",
    "US7376867B1",
    "US20060156312A1",
    "US20070214355A1",
    "US20090260012A1",
    "US20100017644A1",
    "US20100082728A1",
    "US20140189270A1",
    "US9172670B1"
  ]
}

Record 3,404 of 8,000 in Patents full text (MLC-0201). Request the full dataset.