IP Library › Granted Patent US 10,484,326
Granted Patent B2
US 10,484,326 · App. 16/242,940 · Granted Nov 19, 2019

Distributed message queue stream verification

Inventors: Dan Norwood (Mountain View, CA); Neha Narkhede (Mountain View, CA); Anna Povzner (San Jose, CA); Joseph Adler (Palo Alto, CA); Yasuhiro Matsuda (Palo Alto, CA); Jay Kreps (Mountain View, CA)
Assignee: Confluent, Inc.
H04L51/30G06F9/546H04L12/1859H04L12/66H04L43/045H04L43/08H04L51/18H04L51/34H04L41/0806H04L43/106H04L69/16
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,484,326
App. No.
16/242,940
Granted
Nov 19, 2019
Kind
B2
Abstract

A stream verification system for a distributed message queue system with metric collectors on each producer and consumer. A producer time stamp allows correlation of sent and received messages. Verification reports are organized by a message topic. A cumulative checksum allows detection of missing or corrupted messages. Verification messages are used to determine if a zero message report means no messages were sent, or rather that the messages weren't received.

Claims (60)

1. An apparatus for verification of messages in a distributed message queue, comprising:

a central verification analysis system, installed on at least one central verification analysis server, and configured to receive and aggregate verification reports and organize verification reports by a message topic, including producer verification reports from a producer metrics collector module and consumer verification reports from a consumer metrics collector module;

a metric management application, installed on a central management server, that presents to users data related to verification information aggregated by the central verification analysis system;

wherein messages from the producer metrics collector module include a cumulative checksum;

wherein producer verification messages are received at least when no other messages are sent in a particular time period;

wherein the central verification analysis system determines (a) whether any messages were lost or duplicated, (b) the time it takes produced messages to be consumed, and (c) uses a combination of verification messages and cumulative checksums to assess the fidelity of the computations performed in (a); and

a graphical user interface provided by the metric management application allowing a user to interact with and view the calculations completed by the central verification analysis system and request receipt of notifications.

2. The apparatus of claim 1 further comprising an API for providing the graphical user interface.

3. The apparatus of claim 1 wherein the verification reports are encapsulated in messages and sent through a distributed message queue system.

4. The apparatus of claim 3 wherein the producer verification reports contain:

counts of messages sent by a producer;

a count on the number of bytes sent by a producer;

a measurement of message sizes; and

a transport time.

5. The apparatus of claim 3 wherein the producer verification reports contain:

aggregate checksums for the data sent by a producer.

6. The apparatus of claim 5 wherein the checksum is a cumulative checksum.

7. The apparatus of claim 6 wherein the cumulative checksum comprises an aggregation of existing checksums.

8. The apparatus of claim 1 wherein producer verification reports are received for each time window, and a report indicating no messages is received during a time window if no messages were sent.

9. The apparatus of claim 1 wherein the verification reports include counts and checksums on individual topics.

10. The apparatus of claim 1 wherein a time stamp is attached to each message, and wherein metrics are grouped by time stamps.

11. A system for verification of messages in a distributed message queue, comprising:

a central verification analysis system, installed on at least one central verification analysis server, and configured to receive and aggregate verification reports and organize verification reports by a message topic;

a consumer metrics collector module installed on at least one consumer computing device, and configured to collect metrics on messages received from a distributed message queue system and send consumer verification reports to the central verification analysis system;

wherein the consumer verification reports contain (a) counts of messages received from a producer, (b) a count on the number of bytes received from a producer, (c) a measurement of message sizes, and (d) a transport time;

wherein the consumer metrics collector module groups metrics by time stamps;

wherein the consumer metrics collector module detects missing sequence numbers;

wherein the consumer metrics collector module sends a report indicating no messages were received during a time window if no messages were received;

wherein the verification reports include counts and checksums on individual topics;

a metric management application, installed on a central management server, that presents to users analysis of verification information aggregated by the central verification analysis system;

wherein the central verification analysis system determines (a) whether any messages were lost or duplicated, (b) the time it takes produced messages to be consumed, and (c) uses a combination of verification messages and cumulative checksums to assess the fidelity of the computations in performed in (a); and

a graphical user interface provided by the metric management application allowing a user to interact with and view the calculations completed by the central verification analysis system and request receipt of notifications, wherein an API provides the graphical user interface.

12. A system for verification of messages in a distributed message queue, comprising:

a central verification analysis system, installed on at least one central verification analysis server, and configured to receive and aggregate verification reports and organize verification reports by a message topic;

a producer metrics collector module installed on a producer computing device and configured to collect metrics on messages sent to a distributed message queue system and send producer verification reports to the central verification analysis system, wherein the producer metrics collector module attaches a time stamp to each message, wherein the producer metrics collector module attaches a sequence number to each message;

wherein the producer verification reports contain (a) counts of messages sent by a producer, (b) a count on the number of bytes sent by a producer, (c) a measurement of message sizes, and (d) a transport time;

wherein the producer metrics collector module sends producer verification reports for each time window; and

wherein the producer metrics collector module sends a report indicating no messages were sent during a time window if no messages were sent;

wherein the verification reports include counts and checksums on individual topics;

a metric management application, installed on a central management server, that presents to users analysis of verification information aggregated by the central verification analysis system;

wherein the producer metrics collector module is configured to add a cumulative checksum to messages sent out by the producer computer;

wherein the producer metrics collector module is configured to send out verification messages at least when no other messages are sent in a particular time period;

wherein the central verification analysis system determines (a) whether any messages were lost or duplicated, (b) the time it takes produced messages to be consumed, and (c) uses a combination of verification messages and cumulative checksums to assess the fidelity of the computations in performed in (a); and

a graphical user interface provided by the metric management application allowing a user to interact with and view the calculations completed by the central verification analysis system and request receipt of notifications, wherein an API provides the graphical user interface.

13. A method for verification of messages in a distributed message queue, comprising:

collecting metrics in a producer metrics collector module on messages sent by a producer to a distributed message queue system;

receiving producer verification reports from a producer;

receiving consumer verification reports from a consumer;

receiving and aggregating verification reports at a central verification analysis system and organizing verification reports by a message topic;

presenting to users an analysis of verification information aggregated by the central verification analysis system;

wherein the messages from the producer include a cumulative checksum;

wherein producer verification messages are received at least when no other messages are sent in a particular time period; and

wherein the central verification analysis system determines (a) whether any messages were lost or duplicated, (b) the time it takes produced messages to be consumed, and (c) uses a combination of verification messages and cumulative checksums sent via the producer and consumer metrics collectors to assess the fidelity of the computations in performed in (a).

14. The method of claim 13 further comprising encapsulating the verification reports in messages and sending encapsulated reports through a distributed message queue system.

15. The method of claim 13 further comprising aggregating checksums for the data sent by a producer to provide an aggregated checksum, wherein the aggregated checksum is a cumulative checksum.

16. The method of claim 15 wherein the cumulative checksum comprises an aggregation of existing checksums.

17. The method of claim 13 producer verification reports are received for each time window, and a report indicating no messages is received during a time window if no messages were sent.

18. The method of claim 13 wherein the verification reports include counts and checksums on individual topics.

19. The method of claim 13 wherein each message includes a time stamp, and further comprising grouping metrics by time stamps.

20. The method of claim 13 wherein a sequence number is attached to each message, so that missing sequence numbers are detectable.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded May 7, 2026
From: CONFLUENT, INC.
To: INTERNATIONAL BUSINESS MACHINES CORPORATION
Reel/Frame 075569/0163 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 11, 2019
From: NORWOOD, DAN; NARKHEDE, NEHA; POVZNER, ANNA; ADLER, JOSEPH; MATSUDA, YASUHIRO; KREPS, JAY
To: CONFLUENT, INC.
Reel/Frame 050689/0937 →
Continuity (3)
Continuation 15494275 · Apr 21, 2017
Provisional Application 62325936 · Apr 21, 2016
Related Publication 20190149504A1 · May 16, 2019
Cited By (3)
US 12,250,189 US 12,367,068 US 12,417,215