Stream Pipelines concepts
Processing stream data
The Data Fabric streaming architecture, consisting of Streaming Ingestion and Infor Stream Pipelines, provides an end-to-end continuous data flow between the event source and the pipelines that are used for delivering data. This enables a real-time processing of source events to the target system.
When data events are ingested into Data Fabric through the Streaming Ingestion method, Infor Stream Pipelines processes them immediately and continuously without waiting for the events to be stored in Data Lake for durability as data objects. This approach ensures an effective data processing by minimizing the data journey and thus accelerating operations. Consequently, users can extract needed insights from their data in real time.
When source systems publish data events in rapid succession, events may not arrive at Stream Pipelines in the order in which the events occurred. For example, a record deletion event can arrive before an earlier update event for the same record when both events are published within the same second.
The upsert load method handles out-of-order events. The method writes an event to the destination table when the incoming variation value is greater than the stored variation value. The method discards an out-of-order event that has a lower variation value. As a result, an older event cannot overwrite a newer record.
Do not delete records or start downstream processing immediately when a record arrives with a deleted or archived indicator. For example, avoid database triggers that delete records when an insert event occurs. An update event for the same record can arrive after the deletion event. If a process removes the record before the update event is processed, the update event can be rejected or lost.
Design destination-side processes that evaluate the variation value, deleted flag, and archived flag. Do not evaluate each event independently.
A data event is a discrete unit of data that represents a change or update. A data event can take the form of a record, row, message, or any other type of data structure that conveys information. Data events are often used in real-time streaming systems to represent changes or updates to data sources, and they are processed by streaming pipelines to enable a real-time data processing and delivery.
For more information on the Data Fabric ingestion method, see Sending data to Data Lake.
Infor Stream Pipelines have these limitations:
- Currently, Infor Stream Pipelines support only Newline-delimited JSON (NDJSON). Payload records must be in the NDJSON format.
- Events larger than 4.5 MB cannot be processed in Infor Stream Pipelines. If an event's size is larger than 4.5 MB, then the event becomes an error in the Replay Queue.
Processing batched data
Infor Stream Pipelines are designed to process data events in real time. However, Infor Stream Pipelines can also process data that comes through the batch ingestion methods, such as Batch Ingestion API or ION, and data from Data Lake. In the case of batched data, the data is published in a batch format from the data source and must be stored in Data Lake before it can reach Infor Stream Pipelines. This results in a delay between the data event and its delivery to the destination.
Handling deleted and archived records
Stream Pipelines does not filter records based on deleted or archived status. Stream Pipelines does not use deleted or archived flags to remove or modify records in the destination database.
Stream Pipelines deliver all records to the destination database, regardless of deleted or archived status. No DELETE statements or UPDATE statements are issued based on deleted or archived flags.
You must decide how the destination system processes records that contain deleted or archived flags.
Deleted records
When a record is deleted in the source application, Streaming Ingestion publishes a live data change event. Stream Pipelines processes the event in the same way as other records. Stream Pipelines upsert the record to the destination table and set the deleted flag column to true. Stream Pipelines do not issue a DELETE statement against the destination table.
Do not configure destination-side processes, such as database triggers, to delete a row immediately when a record arrives with the deleted flag set to true. An update event for the same record can remain in transit and arrive after the deletion event because events can arrive out of order. Removing the row before the update event is processed can cause the update event to be lost.
Design downstream processes that evaluate the variation value together with the deleted flag. Do not evaluate each event independently.
Archived records
Archived records represent a planned state change in the source application. An archived record is no longer active, but the source application retains the record for historical, legal, or tax purposes.
Archived records typically arrive through Batch Ingestion. Archiving is usually a scheduled or bulk operation, not a continuous stream of individual changes.
Stream Pipelines process archived records in the same way as other incoming records. Stream Pipelines upsert the record to the destination table and set the archived flag column to true. The archived status is stored in the destination table.
Design downstream processes that evaluate the archived flag in the destination table. Do not assume that Stream Pipelines excludes archived records from delivery.
Initial load
Through the Initial Load feature in Infor Stream Pipelines, you can process historical data events from the Data Lake data objects. Initial Load enables you to seed or repopulate destination tables with data that is stored in Data Lake for the entire available history or a selected time window.
All events that are processed by Initial Load contribute to the licensed daily event volume, similarly to live-streamed events. When you plan an initial load, consider the generated additional event volume, especially for large data objects or extended time windows.
See the Infor OS Service Limits documentation.
Initial load duration depends on the data volume and number of record variations that are stored in Data Lake for the selected time frame. The processing time increases when a data object contains many small files rather than fewer large files.
A limited number of initial loads can run concurrently. When multiple initial loads run for pipelines in the same tenant, they are automatically queued by Stream Pipelines to prevent service overload. The queued initial loads are automatically started when preceding loads are completed.
This table shows common initial load use cases:
| Use case | Description |
|---|---|
| Seeding a new destination from a source application |
When a source application is configured with a pipeline for the first time, the destination table is empty. Historical data is published by the source application to Data Fabric through the application's ingestion mechanism. The events flow through the pipeline and populate the destination table. For instructions on how to publish initial data, see the source application's documentation. |
| Seeding or repopulating a destination from Data Lake | Use the Initial Load feature in a pipeline to populate or restore a destination table from data that is stored in Data Lake, for example, after configuring a new destination or recreating a table. Load all available data or limit the load by specifying a time range in the From and To fields. |
| Closing a data gap after pipeline downtime | If a pipeline is stopped, for example, for maintenance or a change in the destination schema, a data gap is created in the destination table for the period in which the pipeline didn't run. To restore missing data, restart the pipeline and run an initial load for the affected time range. Load only the missing period to reduce event volume and processing time. |
Loading data to a destination
Data transfer to a destination through Infor Stream Pipelines is based on the Data Fabric concepts. This includes adherence of data and tables to the schemas that are defined in the Data Catalog and application of data versioning techniques.
Infor Stream Pipelines are designed to handle data loading continuously as events come in. To optimize this process, data is delivered in small batches. Each batch contains data from the time frame of 100 milliseconds or accumulates up to 1000 records, whichever condition is met first.
Stream Pipelines support two load behaviors that determine how events are written to the destination table:
- Upsert
Upsert inserts new records into the destination table and updates existing records.
Stream Pipelines determine whether to update an existing record by comparing the variation value of the incoming event with the variation value that is stored in the destination table. If the incoming event has a higher variation value, Stream Pipelines update the record. If the incoming event has a lower variation value, Stream Pipelines do not write the event to the destination table. This behavior prevents older versions of a record from overwriting newer versions.
Upsert requires a unique identifier for each record in the destination table. Create a primary key or unique index on the destination table before you start the pipeline.
Upsert supports all destination types except Snowflake.
When you upsert records and there are multiple events in the batch window that are associated with the same record, only the events' highest variation is retained for the delivery. Older or equal variations are excluded from the delivery. Events that are excluded during the process are counted and their quantity is displayed on the Excluded graph on the Overview tab in a pipeline.
- Insert
Insert writes every record variation to the destination table as a separate row. Stream Pipelines do not compare variation values when Insert is used. Stream Pipelines append records in the order that the pipeline receives them.
Insert is the supported load method for Snowflake destinations. The Snowflake destination uses Snowpipe Streaming for data transfer instead of a JDBC connection. Snowpipe Streaming does not support update or merge operations.
Do not create a primary key or unique index on the destination table when you use Insert. Primary key and unique index constraints can cause the database to reject events when duplicate records arrive.
Example of a data load with the Upsert method
These conditions must be met to load data with the Upsert method:
- Use the stream pipeline for which the Upsert loading method is defined.
- The Upsert method is supported by the destination.
- Source data has identifier and variation columns.
- The destination table has primary or unique keys columns that match the source identifier columns.
This table shows the destination table that contains these records:
| ID | Description | Price | VariationNumber |
|---|---|---|---|
| 1 | Banana | 1.10 | 9999 |
| 2 | Apple | 0.30 | 9999 |
This table shows events that are published by the source within the time frame of 100 ms:
| ID | Description | Price | VariationNumber |
|---|---|---|---|
| 1 | Banana | 1.50 | 10001 |
| 2 | Apple | 0.50 | 10001 |
| 3 | Orange | 2.10 | 10001 |
| 4 | Watermelon | 3.99 | 10001 |
| 2 | Apple | 0.55 | 10002 |
The upsert optimization process excludes the event with ID 2 and variation number 10001. This is because the delivery batch contains a newer event for the same ID but with a higher value of 10002 in the VariationNumber field.
This table shows records in the destination table after the successful event delivery:
| ID | Description | Price | VariationNumber |
|---|---|---|---|
| 1 | Banana | 1.50 | 10001 |
| 2 | Apple | 0.55 | 10002 |
| 3 | Orange | 2.10 | 10001 |
| 4 | Watermelon | 3.99 | 10001 |
Events with IDs 3 and 4 have been inserted and events with IDs 1 and 2 have been updated.
Example of a data load with the Insert method
These conditions must be met to load data with the Insert method:
- Use the stream pipeline for which the Insert loading method is defined.
- Source data has identifier and variation columns.
- The destination table has no unique column constraints.
This table shows the destination table that contains these records:
| ID | Description | Price | VariationNumber |
|---|---|---|---|
| 1 | Banana | 1.10 | 9999 |
| 2 | Apple | 0.30 | 9999 |
This table shows events that are published by the source:
| ID | Description | Price | VariationNumber |
|---|---|---|---|
| 1 | Banana | 1.50 | 10001 |
| 2 | Apple | 0.50 | 10001 |
| 3 | Orange | 2.10 | 10001 |
| 4 | Watermelon | 3.99 | 10001 |
| 2 | Apple | 0.55 | 10002 |
This table shows records in the destination table after the successful event delivery:
| ID | Description | Price | VariationNumber |
|---|---|---|---|
| 1 | Banana | 1.10 | 9999 |
| 2 | Apple | 0.30 | 9999 |
| 1 | Banana | 1.50 | 10001 |
| 2 | Apple | 0.50 | 10001 |
| 3 | Orange | 2.10 | 10001 |
| 4 | Watermelon | 3.99 | 10001 |
| 2 | Apple | 0.55 | 10002 |
After the successful delivery, the destination table contains all added and already existing events.