Atlas Stream Processing allocates resources per stream processor according to tiers. Fixed resource allocation and cost provides predictability, simplifying the system design process. Use this guide to understand which tiers are most appropriate for your stream processing workloads when planning a deployment.
Resource Allocation
Each tier provides a fixed allocation of processing power, memory, bandwidth, parallelism, and—for processors with Apache Kafka sources—partitions.
Tier | vCPU | RAM (GB) | Bandwidth (Mbps) | Maximum Parallelism | Source Kafka Partition Limit | Initial Sync Collection Limit |
|---|---|---|---|---|---|---|
SP2 | 0.25 | 0.5 | 50 | 1 | 32 | 1 |
SP5 | 0.5 | 1 | 125 | 2 | 64 | 1 |
SP10 | 1 | 2 | 200 | 8 | Unlimited | 5 |
SP30 | 2 | 8 | 750 | 16 | Unlimited | 10 |
SP50 | 8 | 32 | 2500 | 64 | Unlimited | 50 |
Parallelism
Parallelism determines how many threads or concurrent requests a stream processor can use to read, enrich, and write data. You configure parallelism on individual pipeline stages, but Atlas Stream Processing enforces it across the processor as a whole against the maximum for the processor's tier shown in the preceding table.
Stages That Support Parallelism
The following stages accept a parallelism value. Each defaults to 1.
Stage | Effect of Higher Values |
|---|---|
The | |
Increases the maximum number of parallel requests made to the | |
Increases the number of threads across which Atlas Stream Processing distributes write operations, which requires both the stream processor and the cluster it writes to use more computational resources. | |
Increases the number of internal write threads that the sink operator uses, distributing write operations across those threads. For sinks that support the | |
Increases the maximum number of parallel requests made to the external function, which requires more computational resources. |
How Parallelism Works at Runtime
A stream processor moves data through the following sequence:
Source -> Buffer -> Transform -> Buffer -> Sink
For a Apache Kafka source, Atlas Stream Processing starts one consumer thread for each source partition, so a processor can read from many partitions at once. The transform stages in the middle of the pipeline still run on a single thread and process the buffered data in batches.
When a stage with a parallelism value greater than 1 receives a batch, its parallel threads or requests process that batch. The batch then moves on to the next stage.
Document Distribution and Ordering
When a stage's parallelism value is greater than 1, Atlas Stream Processing distributes documents across threads as follows:
If you don't specify
partitionBy, Atlas Stream Processing assigns documents to threads in round-robin order.If you specify
partitionBy, Atlas Stream Processing sends all documents that have the samepartitionByvalue to the same thread, which processes them in order.
$merge is an exception: it uses the fields in its on clause as the partition key.
Cumulative Parallelism
Each stream processor has a maximum cumulative parallelism value determined by its tier. The cumulative parallelism of a stream processor is calculated as follows:
parallelism total - parallelized stages
Where parallelism total is the sum of all parallelism values greater than 1 across the $source, $lookup, $merge, $emit, and $externalFunction stages, and parallelized stages is the number of these stages with parallelism values greater than 1.
For example, if your $source stage sets a parallelism value of 4, your $lookup stage sets no parallelism value (thus defaulting to 1), and your $merge stage sets a parallelism value of 2, then you have two parallelized stages, and the cumulative parallelism of your stream processor is calculated as (4 + 2) - 2.
If a stream processor exceeds the maximum cumulative parallelism for its tier, Atlas Stream Processing throws an error and advises you of the minimum processor tier required for your intended level of parallelism. You must either scale the processor up to a higher tier or lower the parallelism values of your stages to resolve the error. To learn more, see Stream Processing.
Kafka Source Partitions
Processors that read from Apache Kafka are also bound by the source partition limit of their tier. The SP2 tier limits a processor to 32 source partitions and the SP5 tier limits a processor to 64. The SP10 tier and above impose no partition limit.
A processor that exceeds the partition limit of its tier fails, and you must scale it up to support the additional partitions. Because a topic can gain partitions while a processor runs, choose a tier with headroom for the partition growth you anticipate. To learn more about Kafka source behavior, see Atlas Stream Processing Limitations.
Monitor Configured Parallelism
To compare the parallelism you configured against the parallelism available to your processor's tier, use the stats.addedParallelism field that processor statistics return. Atlas Stream Processing returns this field only if at least one stage sets a parallelism value greater than 1.
Workload Selection
The differing resource allocations of each tier make them suitable for different stages and scales of project.
Tier | Use Case |
|---|---|
SP2 | Development, Trial Deployment The lowest-cost option, capable of supporting basic workloads with limited resource requirements. |
SP5 | Development, Basic Production Deployment A low-cost option suitable for production tasks with low throughput, even those employing more complex computation. SP5 processors can support basic filtering, projections, and change stream processing. |
SP10 | Mainstream Production Deployment A baseline for production workloads. SP10 and above are intended for pipelines that require higher levels of parallelism, unlimited Kafka partitioning, or data enrichment operations such as lookups and joins. |
SP30 | Complex Production Deployment A high-performance option designed for memory-intensive stateful operations. SP30 supports pipelines that use long-duration windows, multiple lookups, and stages that require large RAM buffers for data enrichment at scale. |
SP50 | Enterprise-Scale Production The highest-performing option, designed for high-throughput streams and extensive transformation logic. SP50 processors are suitable for operations requiring massive parallelism or compute-intensive workflows. |
Considerations
Consider the following factors when selecting an appropriate tier:
Initialization Surge
A stream processor might require more resources during its initial run than during regular operations. For example, a processor that performs an initialSync against a large Atlas collection needs to support heavy I/O and computation for the duration of the synchronization.
To absorb such elevated demand, select a higher tier temporarily, and scale the processor down when the synchronization is complete and the processor transitions to consuming only new change stream events.
Pipeline Logic
Aggregation pipeline logic is the primary driver of CPU and RAM consumption.
- Windows: Long-lasting windows consume more RAM to hold in-flight
- documents.
- Custom Logic: Javascript
$functionstages or complex grouping - logic increase the computational requirements of each message.
- Custom Logic: Javascript
- Compounding Complexity: Additional stateful or computationally complex stages
- introduce more potential variation in resource demand. Maintaining surplus capacity ensures consistent throughput even during consumption spikes.
Infrastructure
Each point of network or storage contact increases a stream processor's overhead.
- Source or Sink Density: Reading from or writing to parallelized
- sources or sinks—such as Apache Kafka topics with their partitions—increases I/O requirements.
- Data Enrichment:
$lookupand$httpsstages; and operations - against Atlas collections to enrich the data in a stream require network bandwidth and connection pooling.
- Data Enrichment:
- Coordination: In complex deployments orchestrating many sources
- and sinks, stream processors can serve as hubs that route the flow of data between each of these nodes. Such processors benefit from the higher throughput of the SP30 and SP50 tiers.
Conversely, high-throughput stream processing workloads can increase demand on connected resources.
- Impact on Atlas: High-volume, parallelized I/O from a stream
processor can exceed the read or write capacity of source or sink Atlas clusters. This can not only increase latency for the processor, but also bottleneck other workloads dependent on those clusters.
To ensure system-wide performance, scale your Atlas clusters proportionally to the processors with which they interact.
Performance and Latency
Performance goals may require a higher-tier processor even when the processing logic is minimal.
- High Throughput: Higher-tier processors better support streams
- that produce events at a high rate.
Low Latency SLAs: The high parallelism offered by higher-tier processors helps ensure that events don't accumulate in a queue when speed matters. In particular, SP50 processors offer four times the threads as SP30 processors.
- Data Enrichment and Caching: When using
$cachedLookupto enrich - streams with large sets of static or slowly-changing reference data, favor higher-tier processors to provide the necessary RAM for caching.
- Data Enrichment and Caching: When using
- Complex Sinks: Certain sinks involve more expensive
- transformations, transactions, and file management overhead. For processors that interact with these sinks, higher tiers help ensure consistent performance and latency.
Scaling
Atlas Stream Processing scaling is vertical. You can scale a processor manually, or you can enable autoscaling to let Atlas Stream Processing adjust the tier for you.
To scale a processor manually, stop it, select a new tier, and restart it. Atlas Stream Processing checkpoints ensure no data is lost during the transition. Monitor the performance of your processors regularly, and adjust their tiers on the basis of the factors described in this guide.
To let Atlas Stream Processing adjust the tier automatically in response to resource usage, enable vertical autoscaling. Use the factors in this guide to choose the minTier and maxTier that bound the range a processor can scale between.