IP Library › Granted Patent US 12,547,639
Granted Patent B2
US 12,547,639 · App. 18/808,541 · Granted Feb 10, 2026

Enriching event streams with entity data

Inventors: Daniel Joseph Lasky (Chattahoochee Hills, GA); Akihiro Ishikawa (San Francisco, CA); Drew Thompson (San Francisco, CA); Jon Anderson (San Francisco, CA); Udit Mehta (San Francisco, CA)
Assignee: Twilio Inc.
G06F16/254G06F16/27
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,547,639
App. No.
18/808,541
Granted
Feb 10, 2026
Kind
B2
Abstract

System and method for enriching a data stream with enrichment data. The system loads data from one or more customer data warehouses into a storage component using an ingest pipeline; receives, at an enrichment pipeline, an incoming data stream; determines, using the enrichment pipeline, an insertion point within the incoming data stream, the insertion point corresponding to a data object mention; determines, using the enrichment pipeline, enrichment data matching the data object mention, the enrichment data being retrieved from the storage component; augments, via the enrichment pipeline, the incoming data stream with the enrichment data at the determined insertion point to generate an enriched data stream, and transmits the enriched data stream to one or more destinations. The data stream can be an event stream. The enrichment data can be entity data. The system can use a reverse extract/transform/load (ETL) model to enable data ingesting and/or data stream enrichment.

Claims (70)

1 . A system comprising:

one or more computer processors;

one or more computer memories; and

a set of instructions stored in the one or more computer memories, the set of instructions configuring the one or more computer processors to perform operations, the operations comprising:

loading data from one or more customer data warehouses into a storage component using an ingest pipeline, the loading of the data comprising:

creating a source corresponding to a table in a customer data warehouse of the one or more customer data warehouses; and

creating a reverse extract-transform-load (ETL) model associated with the source, the reverse ETL model comprising a relational query comprising table information for the table, the reverse ETL model remaining unchanged upon detecting an update to a schema of the table;

receiving, at an enrichment pipeline, a data stream;

determining, by the enrichment pipeline, an insertion point within the data stream, the insertion point corresponding to a data object mention;

determining, by the enrichment pipeline, enrichment data matching the data object mention, the enrichment data being retrieved from the storage component;

augmenting, using the enrichment pipeline, the data stream with the enrichment data at the determined insertion point to generate an enriched data stream; and

transmitting the enriched data stream to one or more destinations.

2 . The system of claim 1 , wherein the ingest pipeline comprises a scheduler component that determines at least one of a timing or a frequency of data synchronization operations between the one or more customer data warehouses and the storage component.

3 . The system of claim 2 , wherein the ingest pipeline comprises a loader component, the loader component enabled to:

receive, from the scheduler component, synchronization information corresponding to a first job to be executed as part of an data ingest task;

upon receiving the synchronization information associated with the first job, create a second job for a data processing engine based on the first job, the second job being associated with an application programming interface (API) to the storage component; and

execute the second job, the executing of the second job comprising one of at least a data write operation, data retrieval or data deletion operation associated with the storage component.

4 . The system of claim 1 , wherein:

the enrichment pipeline uses a data processing engine associated with an execution plan; and

upon receiving an incoming data stream and detecting that a downstream component is configured to receive an enriched data stream, adding an execution graph node to the execution plan for the data processing engine, the execution graph node associated with a call to an enrichment endpoint for an API to the storage component.

5 . The system of claim 4 , wherein:

the data object mention corresponds to an entity ID; and

determining the insertion point within the incoming data stream comprises detecting the entity ID in the incoming data stream using a path rule and the enrichment endpoint for the API associated with the storage component.

6 . The system of claim 5 , wherein:

the enrichment data comprises one or more entity attributes, each entity attribute associated with at least one attribute value;

the enrichment data matching the data object mention comprises an entity attribute of the one or more entity attributes matching the entity ID based on a matching criterion; and

the enrichment data is retrieved from the storage component using the API.

7 . The system of claim 1 , wherein the data stream corresponds to an event stream and the storage component corresponds to a cache component.

8 . The system of claim 1 , wherein:

the source and the reverse ETL model are associated with an entity model; and

the reverse ETL model further comprises synchronization schedule information associated with the table corresponding to the source.

9 . The system of claim 8 , the operations further comprising creating a mapping between the reverse ETL model and a destination of the one or more destinations.

10 . The system of claim 9 , the operations further comprising:

displaying one of at least the reverse ETL model, the entity model and the mapping between the reverse ETL model and the destination in a user interface (UI);

upon receiving user input indicative of a synchronization failure associated with the source or of a revision to the synchronization schedule information, updating the reverse ETL model; and

upon receiving user input indicative of a revision to the mapping between the reverse ETL model and the destination, updating the mapping.

11 . A method comprising:

loading data from one or more customer data warehouses into a storage component using an ingest pipeline, the loading of the data comprising:

creating a source corresponding to a table in a customer data warehouse of the one or more customer data warehouses; and

creating a reverse extract-transform-load (ETL) model associated with the source, the reverse ETL model comprising a relational query comprising table information for the table, the reverse ETL model remaining unchanged upon detecting an update to a schema of the table;

receiving, at an enrichment pipeline, a data stream;

determining, by the enrichment pipeline, an insertion point within the data stream, the insertion point corresponding to a data object mention;

determining, by the enrichment pipeline, enrichment data matching the data object mention, the enrichment data being retrieved from the storage component;

augmenting, using the enrichment pipeline, the data stream with the enrichment data at the determined insertion point to generate an enriched data stream; and

transmitting the enriched data stream to one or more destinations.

12 . The method of claim 11 , wherein the ingest pipeline comprises a scheduler component that determines at least one of a timing or a frequency of data synchronization operations between the one or more customer data warehouses and the storage component.

13 . The method of claim 12 , wherein the ingest pipeline comprises a loader component, the loader component enabled to:

receive, from the scheduler component, synchronization information corresponding to a first job to be executed as part of an data ingest task;

upon receiving the synchronization information associated with the first job, creating a second job for a data processing engine based on the first job, the second job being associated with an application programming interface (API) to the storage component; and

execute the second job, the executing of the second job comprising one of at least a data write operation, data retrieval or data deletion operation associated with the storage component.

14 . The method of claim 11 , wherein:

the enrichment pipeline uses a data processing engine associated with an execution plan; and

upon receiving an incoming data stream and detecting that a downstream component is configured to receive an enriched data stream, adding an execution graph node to the execution plan for the data processing engine, the execution graph node associated with a call to an enrichment endpoint for an API to the storage component.

15 . The method of claim 14 , wherein the data object mention corresponds to an entity ID, and determining the insertion point within the incoming data stream comprises detecting the entity ID in the incoming data stream using a path rule and the enrichment endpoint for the API associated with the storage component.

16 . The method of claim 15 , wherein:

the enrichment data comprises one or more entity attributes, each entity attribute associated with at least one attribute value;

the enrichment data matching the data object mention comprises an entity attribute of the one or more entity attributes matching the entity ID based on a matching criterion; and

the enrichment data is retrieved from the storage component using the API.

17 . The method of claim 11 , wherein:

the source and the reverse ETL model are associated with an entity model; and

the reverse ETL model further comprises synchronization schedule information associated with the table corresponding to the source.

18 . A non-transitory computer-readable storage medium, the computer-readable storage medium including instructions that when executed by a computer, cause the computer to:

load data from one or more customer data warehouses into a storage component using an ingest pipeline, the loading of the data comprising:

creating a source corresponding to a table in a customer data warehouse of the one or more customer data warehouses; and

creating a reverse extract-transform-load (ETL) model associated with the source, the reverse ETL model comprising a relational query comprising table information for the table, the reverse ETL model remaining unchanged upon detecting an update to a schema of the table:

receive, at an enrichment pipeline, a data stream;

determine, by the enrichment pipeline, an insertion point within the data stream, the insertion point corresponding to a data object mention;

determine, by the enrichment pipeline, enrichment data matching the data object mention, the enrichment data being retrieved from the storage component;

augmenting, using the enrichment pipeline, the data stream with the enrichment data at the determined insertion point to generate an enriched data stream; and

transmit the enriched data stream to one or more destinations.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 20, 2024
From: LASKY, DANIEL JOSEPH; ISHIKAWA, AKIHIRO; THOMPSON, DREW; ANDERSON, JON; MEHTA, UDIT
To: TWILIO INC.; TWILIO INC.
Reel/Frame 069651/0935 →
Continuity (2)
Provisional Application 63534030 · Aug 22, 2023
Related Publication 20250068641A1 · Feb 27, 2025
References Cited (8)
US 11269913B1 · Dervay · 2022 [cited by examiner]
US 20160026692A1 · Cannaliato et al. · 2016 [cited by applicant]
US 20210374143A1 · Neill · 2021 [cited by examiner]
US 20240362196A1 · Gupta · 2024 [cited by examiner]
“International Application Serial No. PCT US2024 042898, International Search Report mailed Nov. 25, 2024”, 4 pgs. [cited by applicant]
“International Application Serial No. PCT US2024 042898, Written Opinion mailed Nov. 25, 2024”, 7 pgs. [cited by applicant]
Barrera, Thalia, “What is Reverse ETL?: Concepts, Use Cases and Integration”, [Online]. Retrieved from the Internet: https: airbyte.com blog reverse-etl, (Aug. 25, 2022). [cited by applicant]
Ehab, Qadah, “An Introduction to Stream Processing with Apache Flink | by Ehab Qadah | Towards Data Science”, (Jan. 7, 2020). [cited by applicant]