IP Library › Granted Patent US 12,541,500
Granted Patent B2
US 12,541,500 · App. 18/667,603 · Granted Feb 3, 2026

Distributed stream-based acid transactions

Inventor: Michael Craig (Wayne, NJ)
Assignee: NASDAQ, INC.
G06F16/2365G06F16/2379
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,541,500
App. No.
18/667,603
Filed
May 17, 2024
Granted
Feb 3, 2026
Kind
B2
Art Unit
2159
USPC
707/703
Abstract

A system for processing distributed transactions is provided. The system includes a sequencer that communicates an atomic message stream to multiple different service instances. The service instances each process the messages from the message stream into a local queue. Each service instance also executes a state machine by reading messages from a queue and transitioning between states in the state machine while also performing one or more operations in connection with performing a distributed transaction.

Claims (96)

1 . A distributed computer system for processing distributed transactions, the distributed computer system comprising:

a plurality of computing devices that communicate by using an electronic data network, each of the plurality of computing devices including at least one hardware processor;

wherein the plurality of computing devices is configured to execute, across different ones of the plurality of computing devices, a sequencer, and a plurality of service instances; and

wherein the sequencer is configured to perform first operations comprising:

receiving a plurality of unsequenced messages;

generating a plurality of sequenced messages that each include a sequence identifier; and

sending, using the electronic data network, the plurality of sequenced messages;

wherein the plurality of service instances are each configured to perform second operations comprising:

executing a state machine that includes a plurality of states that include at least a first state, a second state, and a third state;

in the first state and based on a first message:

attempting to perform a transaction operation for a first distributed transaction, the transaction operation based on contents of the first sequenced message;

transmitting, to the sequencer, an unsequenced result message that includes a vote result of the transaction operation for the first distributed transaction; and

transitioning the state machine from the first state to another one of the plurality of states;

in the second state:

processing second message(s) during a voting period;

closing the voting period based on processing at least one of the second message(s);

based on the closing of the voting period, determining whether to abort or commit the first distributed transaction;

based on a determination to commit the first distributed transaction, committing the transaction operation;

transmitting, to the sequencer, an unsequenced confirmation message that includes a confirmation result that is based on the determining, and an identifier that identifies the service instance on which the state machine is executing; and

transitioning the state machine from the second state to another one of the plurality of states in accordance with transmission of the unsequenced confirmation message;

in the third state:

processing third message(s) until one of the third messages(s) includes an identifier for the service instance on which the state machine is executing; and

as a result of processing the one of the third messages(s), transitioning the state machine to another one of the plurality of states.

2 . The distributed computer system of claim 1 , wherein unsequenced result message includes an operation result from attempting to perform the transaction operation.

3 . The distributed computer system of claim 1 , wherein the state machine consists of three states, the first state being a ready state, the second state being a vote state, and the third state being a confirmation state.

4 . The distributed computer system of claim 3 , wherein the state machine transitions from the ready state to the vote state, from the vote state to the conformation state, and from the confirmation state to the ready state.

5 . The distributed computer system of claim 1 , wherein a plurality of state machine messages includes the first message, the second message(s), and the third message(s), wherein the second operations further comprise:

maintaining a current transaction queue that includes a plurality of sub-queues; and

as part of executing the state machine in each of the first, second and third states, dequeuing the plurality of state machine messages from the current transaction queue.

6 . The distributed computer system of claim 5 , wherein execution of the state machine is performed in a first execution thread and queuing of the plurality of state machine messages to the current transaction queue is performed in a second execution thread that concurrently executes with the first execution thread.

7 . The distributed computer system of claim 5 , wherein the current transaction queue includes a plurality of sub-queues.

8 . The distributed computer system of claim 7 , wherein the second operations further comprise:

dequeuing a message from one of the plurality of sub-queues of the current transaction queue based on which one of the plurality of states the state machine is currently in.

9 . The distributed computer system of claim 7 , with the first message is dequeued from a first sub-queue, the second message(s) are dequeued from a second sub-queue, and the third message(s) are dequeued from a third sub-queue.

10 . The distributed computer system of claim 7 , wherein messages dequeued from one of the sub-queues include both vote message(s) and a timeout message.

11 . The distributed computer system of claim 1 , wherein the voting period is closed as a result of processing one of:

1) a timeout message; and

2) a commit vote message from each identified service participating in the first distributed transaction.

12 . The distributed computer system of claim 11 , wherein the voting period is not closed as a result of processing an abort message.

13 . The distributed computer system of claim 1 , wherein the second operations further comprise:

wherein in the first state and based on the first message, setting a timeout period.

14 . The distributed computer system of claim 13 , wherein the timeout period is set based on a sequence identifier of at least one of the plurality of sequenced messages, wherein closing of the voting period is further based on expiration of the timeout period.

15 . The distributed computer system of claim 14 , wherein the determination of expiration of the timeout period is further based on updating a local time value of the service instance on which the state machine is executing as a result of processing a sequence identifier of another one of the plurality of sequenced messages.

16 . The distributed computer system of claim 1 , wherein the plurality of service instances each include a local database,

wherein the transaction operation is a database operation that is performed against the local database,

wherein committing the transaction operation commits the database operation to the local database.

17 . A method of processing distributed database transactions in a distributed computer system that includes a plurality of computing devices that communicate by using an electronic data network, each of the plurality of computing devices including at least one hardware processor, the method comprising:

executing, across different ones of the plurality of computing devices, a sequencer, and a plurality of service instances;

at the sequencer:

receiving a plurality of unsequenced messages;

generating a plurality of sequenced messages that each include a sequence identifier; and

sending, using the electronic data network, the plurality of sequenced messages;

at each of the plurality of service instances:

executing a state machine that includes a plurality of states that include at least a first state, a second state, and a third state;

in the first state and based on a first message:

attempting to perform a transaction operation for a first distributed transaction, the transaction operation based on contents of the first sequenced message;

transmitting, to the sequencer, an unsequenced result message that includes a vote result of the transaction operation for the first distributed transaction; and

transitioning the state machine from the first state to another one of the plurality of states;

in the second state:

processing second message(s) during a voting period;

closing the voting period based on processing at least one of the second message(s);

based on the closing of the voting period, determining whether to abort or commit the first distributed transaction;

based on a determination to commit the first distributed transaction, committing the transaction operation;

transmitting, to the sequencer, an unsequenced confirmation message that includes a confirmation result that is based on the determining, and an identifier that identifies the service instance on which the state machine is executing; and

transitioning the state machine from the second state to another one of the plurality of states in accordance with transmission of the unsequenced confirmation message;

in the third state:

processing third message(s) until one of the third messages(s) includes an identifier for the service instance on which the state machine is executing; and

as a result of processing the one of the third messages(s), transitioning the state machine to another one of the plurality of states.

18 . The method of claim 17 , wherein a plurality of state machine messages includes the first message, the second message(s), and the third message(s), wherein the method further comprises:

at each of the plurality of service instances, maintaining a current transaction queue that includes a plurality of sub-queues; and

at each of the plurality of service instances, as part of executing the state machine in each of the first, second and third states, dequeuing the plurality of state machine messages from the current transaction queue.

19 . The method of claim 17 , wherein the voting period is closed as a result of processing one of:

1) a timeout message; and

2) a commit vote message from each identified service participating in the first distributed transaction.

20 . A non-transitory computer readable storage medium storing instructions for use with a distributed computer system that includes a plurality of computing devices that communicate by using an electronic data network, each of the plurality of computing devices including at least one hardware processor, the method comprising, the stored instructions comprising instructions that are configured to cause at least one hardware processor to perform operations comprising:

executing, across different ones of the plurality of computing devices, a sequencer, and a plurality of service instances;

at the sequencer:

receiving a plurality of unsequenced messages;

generating a plurality of sequenced messages that each include a sequence identifier; and

sending, using the electronic data network, the plurality of sequenced messages;

at each of the plurality of service instances:

executing a state machine that includes a plurality of states that include at least a first state, a second state, and a third state;

in the first state and based on a first message:

attempting to perform a transaction operation for a first distributed transaction, the transaction operation based on contents of the first sequenced message;

transmitting, to the sequencer, an unsequenced result message that includes a vote result of the transaction operation for the first distributed transaction; and

transitioning the state machine from the first state to another one of the plurality of states;

in the second state:

processing second message(s) during a voting period;

closing the voting period based on processing at least one of the second message(s);

based on the closing of the voting period, determining whether to abort or commit the first distributed transaction;

based on a determination to commit the first distributed transaction, committing the transaction operation;

transmitting, to the sequencer, an unsequenced confirmation message that includes a confirmation result that is based on the determining, and an identifier that identifies the service instance on which the state machine is executing; and

transitioning the state machine from the second state to another one of the plurality of states in accordance with transmission of the unsequenced confirmation message;

in the third state:

processing third message(s) until one of the third messages(s) includes an identifier for the service instance on which the state machine is executing; and

as a result of processing the one of the third messages(s), transitioning the state machine to another one of the plurality of states.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 21, 2024
From: CRAIG, MICHAEL
To: NASDAQ, INC.
Reel/Frame 068737/0001 →
Continuity (1)
Related Publication 20250355860A1 · Nov 20, 2025
References Cited (22)
US 6115365A · Newberg et al. · 2000 [cited by applicant]
US 11178091B1 · Madhavan et al. · 2021 [cited by applicant]
US 11494839B2 · Craig et al. · 2022 [cited by applicant]
US 11503108B1 · Prem et al. · 2022 [cited by applicant]
US 20040240444A1 · Matthews et al. · 2004 [cited by applicant]
US 20100191884A1 · Holenstein et al. · 2010 [cited by applicant]
US 20140040194A1 · Krasensky et al. · 2014 [cited by applicant]
US 20140068635A1 · Holzleitner et al. · 2014 [cited by applicant]
US 20170255668A1 · Schnell · 2017 [cited by examiner]
US 20190384628A1 · Moothoor et al. · 2019 [cited by applicant]
US 20200167865A1 · Craig et al. · 2020 [cited by applicant]
US 20220141826A1 · Ramu et al. · 2022 [cited by applicant]
US 20220150204A1 · Madhavan et al. · 2022 [cited by applicant]
US 20230065589A1 · Gupta · 2023 [cited by examiner]
US 20230174341A1 · Wang et al. · 2023 [cited by applicant]
US 20240152429A1 · Goldstein et al. · 2024 [cited by applicant]
Azure Architecture Center “Saga distributed transactions pattern”, <https://learn.microsoft.com/en-us/azure/architecture/reference-architectures/saga/saga>, seven pages, Sep. 17, 2022. [cited by applicant]
Dremio—Two-Phase Commit <https://www.dremio.com/wiki/two-phase-commit/>, five pages, Oct. 5, 2023. [cited by applicant]
Kleppmann, “Distributed Systems 7.1: Two-phase commit”, Fault-tolerant two-phase commit (1/2): 13:00-18:44, YouTube Video, <https://www.youtube.com/watch?v=-_rdWB9hNlc>, Oct. 28, 2020. [cited by applicant]
Kleppmann, “Distributed Systems Lecture Series”, YouTube Video, <https://www.youtube.com/playlist?list=PLeKd45zvjcDFUEv_ohr_HdUFe97RItdiB>, seven pages, Oct. 28, 2020. [cited by applicant]
Kleppmann, “Distributed Systems” Computer Science Tripos, Part IB, University of Cambridge, <https://www.cl.cam.ac.uk/teaching/2122/ConcDisSys/dist-sys-notes.pdf>, Jan. 2022. [cited by applicant]
Office Action for U.S. Appl. No. 18/667,463, 22 pages, dated Feb. 26, 2025. [cited by applicant]