IP Library › Granted Patent US 12,639,258
Granted Patent B1
US 12,639,258 · App. 18/902,225 · Granted May 26, 2026

Log storage in distributed data streaming systems

Inventors: Vaibhav Sharma (Sammamish, WA); Nagarjuna Koduru (Sammamish, WA); Sayantan Chakravorty (Sammamish, WA); Sai Maddali (Bothell, WA); Usama Bin Naseem (Bothell, WA); Divij Vaidya (Berlin, DE); Mehari Beyene (Everett, WA); Karthikeyan Rajagopalan (Maple Valley, WA)
Assignee: Amazon Technologies, Inc.
G06F16/113
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,639,258
App. No.
18/902,225
Filed
Sep 30, 2024
Granted
May 26, 2026
Kind
B1
Art Unit
2161
USPC
707/672
Abstract

Techniques for log storage in distributed data streaming systems are described. A cluster of brokers receive log records from publishers and send log records to subscribers. The log is represented as a group of segments, each segment subdivided into chunks. Metadata describes the log structure. Log records are stored in chunks at least in a remote storage location shared amongst the brokers in the cluster.

Claims (78)

1 . A computer-implemented method comprising:

by a leader broker in a cluster of brokers, providing data streaming from one or more publishers to one or more subscribers:

generating a chunk object by accumulating records received from the one or more publishers in a local storage location of the leader broker, wherein the chunk object is one of a plurality of chunk objects forming a log that stores records of a data stream;

grouping the plurality of chunk objects into a plurality of segments;

archiving the chunk object to a remote storage location via a network interface of the leader broker after a number of accumulated records or an elapsed time reaches a threshold; and

updating log metadata to include an identification of the archived chunk object; and

by a follower broker in the cluster of brokers:

generating, using the log metadata, a view of locations of records in the log, wherein the view identifies records in chunk objects stored both in the remote storage location and a local storage location of the follower broker, wherein a particular record of the records appears in at least two segments of the plurality of segments and the view identifies the particular record in only one of the at least two segments;

receiving a read request from a subscriber in the one or more subscribers;

identifying, based at least in part on the view of locations of records in the log, the archived chunk object as including at least one record of the data stream responsive to the read request;

reading the at least one record from the remote storage location via a network interface of the follower broker; and

sending the at least one record to the subscriber.

2 . The computer-implemented method of claim 1 , wherein the cluster of brokers is managed by a data streaming service of a cloud provider network, and further comprising:

launching a new broker to add to the cluster of brokers; and

by the new broker, sending a request for records to the leader broker, the request for records in chunk objects that have yet to be archived to the remote storage location.

3 . The computer-implemented method of claim 1 , wherein the cluster of brokers includes at least one read-only broker, and further comprising, by the leader broker, sending records received from the one or more publishers to the follower broker and not to the read-only broker.

4 . A computer-implemented method comprising:

writing, by a first broker in a cluster of brokers implemented by one or more physical computing devices, a chunk object to a remote storage location via a network interface, the chunk object including records received from one or more publishers of a data stream, wherein the chunk object is one of a plurality of chunk objects forming a log structure that stores the records of the data stream;

grouping the plurality of chunk objects into a plurality of segments;

updating log metadata to include an identification of the chunk object;

generating, by a second broker in the cluster of brokers, using the log metadata, a view of locations of records in the log, wherein the view identifies records in chunk objects stored both in the remote storage location and a local storage location of the second broker, and wherein a particular record of the records appears in at least two segments of the plurality of segments and the view identifies the particular record in only one of the at least two segments;

receiving, by the first broker, a read request from a subscriber of the data stream;

identifying, by the first broker, based at least in part on the view of locations of records in the log, the chunk object as including at least one record of the data stream responsive to the read request;

reading, by the first broker, the at least one record from chunk object in the remote storage location via the network interface; and

sending, by the first broker, the at least one record to the subscriber.

5 . The computer-implemented method of claim 4 , further comprising:

obtaining, by the second broker, the log metadata including a log manifest that identifies the log structure as one or more segments and, for each segment, a segment manifest that identifies the plurality of chunk objects included in that segment;

updating, by the second broker using the log manifest and the one or more segment manifests, the view of locations of records in the log;

receiving, by the second broker, a second read request from the subscriber of the data stream;

identifying, by the second broker based at least in part on the view of locations of records in the log, a second chunk object as including at least a second record of the data stream responsive to the second read request;

reading, by the second broker, the at least the second record from the second chunk object in the remote storage location via the network interface; and

sending, by the second broker, the at least the second record to the subscriber.

6 . The computer-implemented method of claim 4 , further comprising, by the first broker, generating the chunk object by accumulating records from the one or more publishers in the local storage location prior to writing the chunk object in the remote storage location.

7 . The computer-implemented method of claim 6 , further comprising:

receiving, by the first broker, a request from a new broker in the cluster of brokers for accumulated records that have yet to be written to the remote storage location; and

sending, by the first broker, the accumulated records that have yet to be written to the remote storage location to the new broker.

8 . The computer-implemented method of claim 4 , further comprising, by the first broker, sending a metadata update to the second broker in the cluster of brokers, the metadata update including an identification of the chunk object written to the remote storage location.

9 . The computer-implemented method of claim 8 , wherein the plurality of chunk objects are grouped into a plurality of log segments, and wherein the log structure includes a log manifest that includes metadata identifying each of the log segments and a segment manifest for each of the plurality of log segments, each segment manifest identifying the chunk objects grouped into the segment, and further comprising:

updating, by the first broker, the segment manifest when a new chunk object is added to a log segment; and

updating, by the first broker, the log manifest when a new log segment is added to the log structure,

wherein the metadata update includes an indication of an update to either the segment manifest or the log manifest.

10 . The computer-implemented method of claim 8 , wherein the metadata update is sent in response to:

reaching a threshold number of chunk objects written to the remote storage location since a previous metadata update; or

reaching a threshold time since a previous metadata update.

11 . The computer-implemented method of claim 4 , wherein the read request includes an identifier of requested data, the computer-implemented method further comprising determining, by the first broker, a location within the chunk object in the remote storage location from which to begin reading data based on a comparison of the identifier to a record index stored in the chunk object.

12 . The computer-implemented method of claim 11 , wherein the identifier is a timestamp of a record in the requested data and the record index is a timestamp index.

13 . The computer-implemented method of claim 11 , wherein the identifier is a record offset of a record and the record index is an offset index.

14 . A system comprising:

a first one or more physical computing devices implementing a data storage service in a multi-tenant provider network, the data storage service providing a remote storage location; and

a second one or more physical computing devices implementing a first broker in a cluster of brokers in the multi-tenant provider network, the first broker including instructions that upon execution cause the first broker to:

write a chunk object to the remote storage location via a network interface of the second one or more physical computing devices, the chunk object including records received from one or more publishers of a data stream, wherein the chunk object is one of a plurality of chunk objects forming a log structure that stores the records of the data stream;

group the plurality of chunk objects into a plurality of segments;

update log metadata to include an identification of the chunk object;

generate, by a second broker in the cluster of brokers, using the log metadata, a view of locations of records in the log, wherein the view identifies records in chunk objects stored both in the remote storage location and a local storage location of the second broker, and wherein a particular record of the records appears in at least two segments of the plurality of segments and the view identifies the particular record in only one of the at least two segments;

receive a read request from a subscriber of the data stream;

identify, based at least in part on the view of locations of records in the log, the chunk object as including at least one record of the data stream responsive to the read request;

read the at least one record from chunk object in the remote storage location via the network interface; and

send the at least one record to the subscriber.

15 . The system of claim 14 , further comprising a third one or more physical computing devices implementing the second broker in the cluster of brokers, the second broker including instructions that upon execution cause the second broker to:

obtain the log metadata including a log manifest that identifies the log structure as one or more segments and, for each segment, a segment manifest that identifies a plurality of chunk objects included in that segment;

update, using the log manifest and the one or more segment manifests, the view of locations of records in the log;

receive a second read request from the subscriber of the data stream;

identify, based at least in part on the view of locations of records in the log, a second chunk object as including at least a second record of the data stream responsive to the second read request;

read the at least the second record from the second chunk object in the remote storage location via the network interface; and

send the at least the second record to the subscriber.

16 . The system of claim 14 , wherein the first broker includes further instructions that upon execution cause the first broker to generate the chunk object by accumulating records from the one or more publishers in the local storage location of the second one or more physical computing devices prior to writing the chunk object in the remote storage location.

17 . The system of claim 16 , wherein the first broker includes further instructions that upon execution cause the first broker to:

receive a request from a new broker in the cluster of brokers for accumulated records that have yet to be written to the remote storage location; and

send the accumulated records that have yet to be written to the remote storage location to the new broker.

18 . The system of claim 14 , wherein the first broker includes further instructions that upon execution cause the first broker to:

send a metadata update to the second broker in the cluster of brokers, the metadata update including an identification of the chunk object written to the remote storage location.

19 . The system of claim 18 , wherein the plurality of chunk objects are grouped into a plurality of log segments, and wherein the log structure includes a log manifest that includes metadata identifying each of the log segments and a segment manifest for each of the plurality of log segments, each segment manifest identifying the chunk objects grouped into the segment, wherein the first broker includes further instructions that upon execution cause the first broker to:

update the segment manifest when a new chunk object is added to a log segment; and

update the log manifest when a new log segment is added to the log structure,

wherein the metadata update includes an indication of an update to either the segment manifest or the log manifest.

20 . The system of claim 18 , wherein the metadata update is sent in response to:

reaching a threshold number of chunk objects written to the remote storage location since a previous metadata update; or

reaching a threshold time since a previous metadata update.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 25, 2026
From: SHARMA, VAIBHAV; KODURU, NAGARJUNA; CHAKRAVORTY, SAYANTAN; MADDALI, SAI; NASEEM, USAMA BIN; VAIDYA, DIVIJ; BEYENE, MEHARI; RAJAGOPALAN, KARTHIKEYAN
To: AMAZON TECHNOLOGIES, INC.
Reel/Frame 074474/0448 →
References Cited (25)
US 6012067A · Sarkar · 2000 [cited by examiner]
US 9575978B2 · Tevis · 2017 [cited by examiner]
US 9729653B2 · Nampally · 2017 [cited by examiner]
US 10037337B1 · Shanmuganathan · 2018 [cited by examiner]
US 10976949B1 · Calhoun, Jr. · 2021 [cited by examiner]
US 11042504B2 · Kashi Visvanathan · 2021 [cited by examiner]
US 11582261B2 · Vivekanandan · 2023 [cited by examiner]
US 12348593B2 · Chakravorty · 2025 [cited by examiner]
US 20050091253A1 · Cragun · 2005 [cited by examiner]
US 20120239623A1 · McCann · 2012 [cited by examiner]
US 20160307274A1 · Sweeney · 2016 [cited by examiner]
US 20180091586A1 · Auradkar · 2018 [cited by examiner]
US 20200097578A1 · Diaconu · 2020 [cited by examiner]
US 20200195572A1 · Efimov · 2020 [cited by examiner]
US 20210096955A1 · Ajith · 2021 [cited by examiner]
US 20230079486A1 · Yarlagadda · 2023 [cited by examiner]
US 20240371510A1 · Aman · 2024 [cited by examiner]
JP 7159388B2 · 2022 [cited by examiner]
“Architecture”; WarpStream; downloaded from <https://docs.warpstream/overview/architecture> on Jul. 31, 2024, 6 pages. [cited by applicant]
“Client Configuration for Bufstream”; Buf Docs; downloaded from <https://buf.build/docs/bufstream/kafka-compatibility/configure-clients#connecting-to-bufstream> on Jul. 31, 2024, 4 pages. [cited by applicant]
“Difference with Apache Kafka”; downloaded from <https://docs.automq.com/automq/what-is-automq/difference-with-apache-kafka> on Jul. 31, 2024, 5 pages. [cited by applicant]
“Difference with Tiered Storage”; downloaded from <https://docs.automq.com/automq/what-is-automq/difference-with-tiered-storage> on Jul. 31, 2024, 3 pages. [cited by applicant]
Kumar, Abhijeet; “KIP-1023: Follower Fetch From Tiered Offset”; Apache Software Foundation; downloaded from <https://cwiki.apache.org/confluence/display/KAFKA/KIP-1023%3A+Follower+fetch+from+tiered+offset> on Jul. 31, 2… [cited by applicant]
“Overview”; AutoMQ; downloaded from <https://docs.automq.com/automq/what-is-automq/overview> on Jul. 31, 2024, 3 pages. [cited by applicant]
Selwan, Marc; “Introducing Confluent Cloud Freight Clusters”; downloaded from <https://www.confluent.io/blog/introducing-confluent-cloud-freight-clusters/> on Jul. 31, 2024, 10 pages. [cited by applicant]