Buffering
Backpressure and the need to buffer
While enterprises typically aim to ensure that their Kron TP deployments are appropriately sized for the expected load, issues can arise. This may occur when an application starts generating more logs than usual, or when the downstream service, where data is sent, begins responding slower than usual.
In the design of Kron TP's topology, part of the strategy involves the transmission of backpressure signals. Backpressure serves as an indication that events cannot be processed as quickly as they are being received. When a component attempts to send more events to another component than it can currently handle, the sending component is indirectly notified. Backpressure can traverse from a sink, through any transforms, back to the source, and ultimately reach clients, such as applications sending logs via HTTP.
Backpressure serves as a mechanism for a system to communicate whether it can handle additional work or if it is too occupied to do so. It allows the system to slow down the consumption or acceptance of events, for example, when pulling them from a source like Kafka or receiving them over a socket like HTTP.
However, there are cases where immediate propagation of backpressure is not always desirable. This is because it could lead to a constant slowdown of upstream components and callers, potentially causing issues beyond the scope of Kron TP. The goal is to avoid a complete slowdown when a component just surpasses the threshold of being fully saturated, and to be able to handle temporary slowdowns and outages in external services to which sinks send data.
Buffering is the approach that Kron TP employs to address these challenges.
Buffering between components
Each component within a Kron TP topology is equipped with a modest in-memory buffer. While the main function of this buffer is to serve as the communication channel between two components, we extend its utility by guaranteeing a small amount of space – usually around 100 events. This provision allows events to be transmitted even when the receiving component is presently occupied, optimizing throughput in scenarios where workloads are not consistently uniform.
Nevertheless, to safeguard against transient overloads or outages, it is imperative to implement a more extensive buffering solution that can be customized to suit the specific workload at hand.
Buffering at the sink
When configuring a Kron TP setup, your focus will be on adjusting buffer configuration settings for sinks. This is primarily because, in practical scenarios, sinks serve as the primary generator of backpressure in a topology, especially when communicating with services over the network, where latency might be introduced, or temporary outages may occur.
By default, sinks utilize an in-memory buffer, much like all other components. However, the default buffer size is slightly augmented to 500 events. This augmentation is specifically applied to sinks since they typically act as the primary source of backpressure in any given Kron TP topology. Apart from the increased default buffer capacity, you also have complete control over the buffer configuration. Kron TP provides access to two key settings for buffer control: the type of buffer to employ and the action to be taken when the buffer reaches its full capacity.
Buffer types
In-memory buffers
We've previously discussed in-memory buffers. As the name implies, this type of buffer stores events in memory. In-memory buffers are the fastest among buffer types, but they have two primary drawbacks: they can consume memory proportionate to their size, and they lack durability.
Despite Kron TP offering features such as end-to-end acknowledgements to ensure that sources refrain from acknowledging events until they have been processed, any events residing in an in-memory buffer would be lost in the event of a Kron TP or host crash. While pull-based sources like S3 or Kafka could handle this by attempting to reprocess the events, push-based sources might not have the capability to retransmit their messages.
Disk buffers
When prioritizing the durability of data over the overall performance of Kron TP, the use of disk buffers becomes essential. Disk buffers enable the persistence of buffered events even during the restart of Kron TP, including incidents of Kron TP crashes. This functionality allows Kron TP to resume operations from where it left off upon restarting.
Operating similarly to a write-ahead log, disk buffers ensure that each event passes through the buffer, is written to the data files, and then read back out. While this process might seem slow in theory, modern operating systems facilitate reads directly from memory, resulting in disk buffers maintaining high throughput on both the read and write paths. By default, data synchronization to disk does not occur with every write but is instead synchronized at intervals (e.g., every 500 milliseconds). This approach balances high throughput with a reduced risk of data loss.
The design of disk buffers is geared toward delivering consistent performance. Although other projects may boast faster data writing to disk, Kron TP prioritizes ensuring that events can be read as swiftly as they are written, minimizing the latency between writing and reading events.
Similar to in-memory buffers, disk buffers have a configurable maximum size to limit disk usage. This maximum size is strictly enforced, ensuring that Kron TP does not exceed it. However, there is a minimum size for all buffers, currently set at approximately 256MiB, as mandated by the disk buffer implementation. On the filesystem, disk buffers appear as append-only log files that grow to a maximum size of 128MiB and are deleted once fully processed.
Recognizing the potential for storage errors, whether due to hardware failures or accidental deletion of data files during Kron TP operation, disk buffers automatically checksum all events written to disk. In the event of corruption detected during a read, disk buffers automatically recover as many events as can be accurately decoded. Metrics are emitted by disk buffers when corruption is detected, providing an accurate insight into the number of lost events.
“When full” behavior
Choosing the action to take when a buffer is full is just as crucial as selecting the buffer type. This decision can significantly influence the overall performance of Kron TP as a system, and it often needs to align with the specific configuration and workload.
Blocking
When set to the "block" configuration, Kron TP will indefinitely pause attempts to write to a full buffer. This represents the default behavior when the buffer reaches capacity.
This default behavior is in place because it generally ensures the desired outcome of reliably processing observability data in the order it was received by Kron TP. Moreover, the blocking mechanism introduces backpressure, which, as discussed earlier, is a crucial signal to upstream components, indicating that they may need to reduce their pace or shed load.
However, blocking might not be suitable in cases where you are receiving data from clients and cannot afford to have them waiting for a response indicating that the data has been accepted by Kron TP. Further details on common buffering scenarios and configurations will be discussed later on.
Drop the event
When set to "drop newest," Kron TP will straightforwardly discard an event if the buffer is presently full.
This functionality proves beneficial when the data itself is idempotent (where the same value is continually sent) or is generally of lower significance, such as trace or debug logging. This approach enables Kron TP to efficiently manage load shedding by reducing the number of events in-flight for a topology, all the while preventing the blocking of upstream components.
Recommended buffering configurations
Here are a few common situations frequently encountered by Kron TP users along with the recommended buffering configurations:
No Storage Provided to Kron TP:
Use in-memory buffers in this case. Kron TP does not support buffering events to external storage systems.
Performance as the Top Priority:
Opt for in-memory buffers. While the drop_newest mode offers the highest performance, be aware that more events may be dropped than expected.
Emphasizing Durability:
Utilize disk buffers for enhanced durability. Depending on your sources, you may choose to maintain the default blocking behavior or opt to drop events when the buffer is full. For sources receiving data directly from clients, dropping the event might be preferable to avoid client waiting, potentially causing issues further up the stack. In general, increasing the max_events and retaining the default blocking behavior is often sufficient to handle higher event processing rates.