How we built a production architecture handling 1,900+ messages per second and tens of thousands of concurrent connections with Kafka, Redis Pub/Sub, .NET Channels, and SSE.

Mofid Brokerage is the largest brokerage in Iran.

I wrote this article based on our team’s real-world experience, where we faced technical challenges and worked closely together to find the right solutions. Our goal was to display real-time stock and fund prices in the Mofid application so our users could have a smooth and up-to-date experience.

That decision marked the beginning of a challenging technical journey. We were dealing with a huge amount of data that had to travel from the trading core to the user’s screen with the lowest possible latency.

At the beginning, there were many uncertainties: Where exactly should we start? How should we identify and remove bottlenecks? And how should we design an architecture that would remain stable during peak market traffic?

On the surface, the requirement sounded simple. In practice, however, we faced challenges that could have pushed the entire system into a dead end if we had not identified them correctly from the beginning.

  • Where exactly does the data come from?
  • At what rate is it generated?
  • How large can it become?
  • More importantly, how can we deliver this amount of data to millions of users without sacrificing latency, stability, and simplicity?
  • What method should we use to send the data to clients?

In this article, I want to share our team’s experience dealing with these challenges, the architectural decisions we made, and the solutions we implemented to achieve this goal.

The Real-Time Price Data Source and Initial Challenges

The first step was gaining a deep understanding of the data source.

Where are real-time prices generated, and at what rate?

Critical market data, including real-time stock and fund prices, is generated by the central core of the Securities and Exchange Organization, RLC.

To receive this data as quickly as possible, another specialized team inside Mofid is responsible for reading this data stream.

They use C++ for this task, which is a smart choice because of its proximity to the hardware and high performance, minimizing latency while reading the information.

After initial processing, the team puts the incoming messages into a Kafka cluster using the optimized Protobuf format to reduce message size and increase serialization speed.

At this point, we had a very fast pipeline transferring data from the exchange core to Kafka.

The next question was:

How much data were we actually dealing with?

The message production rate is directly related to the market’s daily trading volume.

But to design the system, we needed an accurate estimate of both normal and peak conditions.

Our analysis showed that we were generally dealing with four main data streams:

  1. Last-trade price queue: around 1,000 messages per second.
  2. Closing-price queue: around 500 messages per second.
  3. User asset-change queue: around 200 messages per second.
  4. User return-change queue: around 200 messages per second.

With a quick calculation, under normal market conditions we had to process approximately 1,900 messages per second.

The important point was that during highly volatile market days, this number could easily become two to three times higher for short periods.

Our architecture had to be ready for the worst-case scenario.

Strategy for Moving Data to Internal Services

The main challenge at this stage was moving this huge volume of messages from Kafka into our backend service layer with the lowest possible latency.

We needed a tool that could act as a very fast intermediary.

We had several options:

  • RabbitMQ Streams
  • Kafka
  • Redis Pub/Sub

Our decision criteria included throughput, latency, scalability, ease of maintenance, and resource consumption.

I had previously worked with RabbitMQ in other projects and knew it had good performance, but we wanted to make a complete comparison.

Our choice: Redis Pub/Sub

Kafka was an excellent choice for ingestion and maintaining the main data stream.

But for the real-time fan-out layer to our SSE services, what we needed was fast, simple, ephemeral broadcasting rather than replay and durability.

For this reason, at this point in the architecture we chose Redis Pub/Sub as a lightweight and low-latency layer between Kafka and the backend.

Our goal at this layer was not to maintain message history or replay messages.

Instead, the goal was to deliver the latest price changes to the backend with the lowest possible latency.

Redis Pub/Sub uses at-most-once semantics, which matched the nature of our data very well.

If a price update from a few seconds ago is lost, it will be replaced by the next price update.

In contrast, Kafka, Redis Streams, and RabbitMQ Streams are more suitable for scenarios where replay, durability, consumer offsets, acknowledgments, or more reliable processing are required.

We already had Redis in our technology stack, so choosing it did not introduce the complexity of maintaining another technology for our SRE team.

Development was also very simple.

Implementation: Working with our SRE team, we used Redpanda Connect to transfer data from Kafka topics into Redis Pub/Sub channels.

One advantage of Redpanda Connect for us was that it allowed us to keep the pipeline between Kafka and Redis declarative and observable without writing a dedicated intermediary service.

To maximize performance, we configured a dedicated Redis node in a completely memory-only setup.

A note about reliability:

In this scenario, if Redis restarted, a few moments of data could be lost and prices could stop updating for a few seconds.

Considering that the main database was updated in the background with approximately five minutes of delay, this level of incident risk was acceptable for displaying real-time prices.

Here is what Redis resource usage looked like in production:

It seems this workload was almost a joke for Redis!

Delivering Data to the Client

At this point, the data was reaching Redis at the speed of light.

The next major challenge was sending these 1,900+ messages per second to hundreds of thousands of connected clients across browsers and mobile devices.

We evaluated several options for real-time communication with clients:

  • SignalR
  • Socket.IO
  • Lightstreamer
  • SSE

We needed a technology that:

  1. Had high scalability.
  2. Worked easily across browsers and devices.
  3. Preferably did not require a heavy client-side library.
  4. Was compatible with our .NET stack.
  5. Was quick and easy to develop.

Our choice: SSE (Server-Sent Events)

Choosing between these options was difficult, but we had clear criteria:

It had to be free, simple to scale, easy to maintain, unaffected by sanctions, fully compatible with our .NET stack, and easy to run across devices.

.NET 10 introduced some classes and capabilities that make working with SSE easier, but that does not mean SSE cannot be used with older versions.

For our requirements, SSE was the clear winner.

Unlike WebSockets, which provide a more complex bidirectional connection, we only needed to push data from the server to the client.

SSE is a very simple web standard that works over HTTP.

The native JavaScript EventSource API is built into browsers, SSE works well with standard HTTP infrastructure, and the protocol has built-in reconnection behavior.

Browser compatibility is also excellent, and SSE has been supported for many years.

We built a Proof of Concept (POC) with GitHub Copilot.

The results were impressive.

The simplicity of the implementation and its excellent performance made our decision clear.

Backend Code Example

[HttpGet]
public async Task GetPrices()
{
Response.Headers.Append("Content-Type", "text/event-stream");
Response.Headers.Append("Cache-Control", "no-cache");
Response.Headers.Append("Connection", "keep-alive");
Response.Headers.Append("X-Accel-Buffering", "no");

var id = Guid.NewGuid().ToString("N");
var cancellationToken = HttpContext.RequestAborted;
var reader = SubscribersCoordinatorHostedService.Subscribe(id, cancellationToken);

await foreach (var (eventName, eventData) in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
{
await Response.WriteAsync($"event: {eventName}\ndata: {eventData}\n\n", cancellationToken).ConfigureAwait(false);
await Response.Body.FlushAsync(cancellationToken).ConfigureAwait(false);
}
}

Frontend Code Example

const source = new EventSource("/prices/stream");

source.onmessage = (event) => {
const price = JSON.parse(event.data);
render(price);
};

source.onerror = () => {
console.log("Retrying...");
};

Client-Side Implementation: Overcoming the Limitations

Implementing the SSE mechanism on the client side was no less challenging than the backend.

We focused on bringing this continuous flow of data to the screen without reducing the quality of the user experience.

On the client side, we faced six major challenges, and we defined a specific strategy for each one.

1. Browser Concurrent Connection Limits

With HTTP/1.1, there is a limit of six simultaneous open connections per domain.

Solution: Fortunately, because our infrastructure was fully running on HTTP/2, this limitation was effectively removed through multiplexing and support for many concurrent streams.

As a result, we got past this bottleneck without requiring special changes on the client or implementing a SharedWorker.

2. Authentication and Sending Tokens

The native JavaScript EventSource class only supports GET and does not provide a simple way to set custom headers such as:

Authorization: Bearer

So how were we supposed to authenticate the user?

We had three options.

Option one: Use cookies and credentials so the browser sends authentication cookies with the request.

Option two: Send the token as a query string in the URL.

From a security perspective, this is not recommended because tokens may be stored in server or proxy logs and could leak.

Option three — our choice: Use the eventsource package.

The package provides an EventSource-compatible API while giving us additional control over the underlying Fetch API.

This allowed us to customize the request, including adding authentication headers such as Authorization, while keeping the familiar EventSource programming model.

We wanted to reduce implementation complexity and keep our main focus on the business logic.

This package helped us do exactly that.

3. The Update Storm

Given the high message production rate, if we updated the page state directly for every event, the browser would have to deal with continuous and expensive rendering.

That would increase the number of renders and reduce UI responsiveness.

Solution: To update multiple prices on the screen at the same time, we used a Batch Update strategy.

We temporarily collected incoming data inside a ref, and then at defined intervals injected all the new data into the UI at once.

This made the user experience significantly smoother and eliminated lag.

4. Error Handling and Reconnection

In the real world of mobile internet and Wi-Fi, connection drops are unavoidable.

Solution: SSE already has a built-in reconnection model.

But for more precise control, we decided to manage the retry behavior ourselves.

When an error occurred, we closed the connection and re-established it after a reasonable delay to prevent unnecessary pressure on the server.

5. Managing Resources While Navigating Between Pages

When users move around the application and open different stock or fund detail pages, connections opened for previous pages are no longer useful.

Keeping those connections alive only wastes client and server resources.

Solution: Cleanup.

When the page or selected fund changes and the component is unmounted, we immediately close the current connection.

Then we open a new connection for the data required by the new page.

This apparently simple step had a major impact on optimizing browser memory usage and reducing server load.

6. Initial Values

The SSE connection was responsible only for real-time data.

To prevent inconsistencies, the client first retrieved a snapshot of the latest prices through the normal API, which had a small delay.

The stream then applied only subsequent changes.

The same snapshot/re-sync process happened after reconnecting so that losing Redis Pub/Sub messages or temporarily losing the internet would not leave stale data on the screen.

Backend Bottleneck — Processing at Maximum Speed

This was where we reached one of the biggest technical challenges.

We had managed to get the data into Redis and had already chosen how to deliver it to the client.

But how were we supposed to read this massive amount of data — approximately 1,900 messages per second from several different queues — inside our .NET service and efficiently distribute it among users?

We were dealing with a combination of public data, such as stock prices, and private data, such as a user’s assets.

Public data was the same for every user.

For example, the price of a stock.

Private data belonged to a particular user.

For example, if a user sold part of their assets or purchased more, they needed to see that change reflected immediately in their portfolio.

Our architecture had to support both types of data.

It also had to remain flexible enough to support new queues in the future.

The challenges were:

  • A high volume of real-time messages.
  • A combination of public and user-specific private data.
  • High concurrency requirements.
  • An extensible architecture for adding new queues in the future.
  • A large number of concurrent users requiring broadcasts.

For this critical part of the .NET implementation, we evaluated two main approaches:

  • TPL Dataflow
  • Channels

Our choice: System.Threading.Channels

Our main criteria were implementation simplicity, development speed, minimum latency, maximum efficiency, and optimized RAM and CPU usage.

TPL Dataflow is a powerful option for more complex pipelines, multi-stage transformations, batching, and actor-like models.

But System.Threading.Channels is designed specifically for high-performance producer/consumer scenarios with low overhead.

Channels in .NET behave like a smart, fully asynchronous FIFO pipe.

In our own experience, they performed extremely well.

Channels are thread-safe and highly optimized for low-overhead producer/consumer scenarios.

Microsoft’s benchmarks also show that in appropriate scenarios they can provide very high throughput with extremely low allocations.

This gave us cleaner and more scalable code without requiring us to manually manage much of the synchronization and queue coordination ourselves.

According to Microsoft’s benchmark, Channels performed significantly better than Dataflow in both speed and memory usage in the scenario being tested.

We designed a layered and modular architecture.

Each Redis Channel was handled by a HostedService, which pushed its messages into an internal Channel.

To prevent queues from growing forever, we used a BoundedChannel with a capacity of 100k and DropOldest.

This means that the Channel has a maximum capacity of 100k.

We did this for two reasons.

First, we did not want an unbounded queue.

If processing became slow, we did not want to keep accumulating messages indefinitely.

Second, having a predefined capacity can improve the behavior and predictability of the Channel under load.

If more messages arrive after the Channel is full, DropOldest removes the oldest message.

We also introduced a Coordinator to manage the Redis Channels.

The Coordinator manages the Channels and creates a separate Channel for each user.

We did not blindly send every message to every user.

Each connection received only the subset of data it actually needed based on the page context, the user’s assets, or their watchlist.

This filtering kept the real fan-out under control and prevented the number of outgoing events from being multiplied by the total number of market symbols.

The architecture is scalable and ready for new queues.

RLC → C++ Ingestion → Kafka/Protobuf → Redpanda Connect → Redis Pub/Sub → .NET Hosted Services → Coordinator → Per-user/per-subscription channels → SSE → Client batch renderer

Go Ahead, Don’t Be Afraid

Once we had developed the architecture and passed our heavy tests across different environments with the help of the SRE team, it was time for the stressful moment of going live.

Our strategy was:

“Go live, don’t be afraid” — but with a seat belt on.

We used Unleash as our Feature Toggle tool.

In the first step, we enabled the feature for only 30% of users.

At the same time, our eyes were on the monitors and Grafana dashboards.

Fortunately, everything looked great.

We saw no problems or performance degradation in either the backend or the infrastructure.

But during that same 30% rollout, we received important feedback from users and our product manager.

The visual effect when prices changed on the screen was not attractive enough and did not communicate the feeling of real-time updates very well.

That was when one of our frontend developers and I started working on it.

We decided to use the price-change effect from TradingView as our reference.

It is a well-known and widely tested effect that provides a great user experience.

The interesting technical part was that we did not want to write additional JavaScript for this visual effect.

Adding JavaScript on the client for this type of processing could introduce additional overhead at scale.

So our frontend team reproduced the TradingView-style effect using a few simple and clean CSS tricks.

It was extremely lightweight and the result looked excellent.

After making those visual changes and confirming the system was stable, we confidently enabled the feature for 50% of users.

We stayed at that level for a few days so we could collect enough data about system and user behavior.

Once we were sure everything was working as expected, we opened the valve a little more:

First 70%, and finally 100%.

The Numbers Speak for Themselves

Designing an architecture on paper is one thing.

Seeing how it behaves under real production load is something else entirely.

One of our biggest concerns was server resource consumption.

Because we expected thousands of simultaneous open connections, memory management was critical.

One of our biggest fears was:

“What if the queues are filling up and we don’t even know it?”

The final result surprised even us.

Thanks to SSE’s lightweight structure and efficient memory management in .NET — particularly System.Threading.Channels, which has very low overhead — we were able to handle a huge amount of traffic with very limited resources.

We used Prometheus to collect metrics and Grafana for visualization.

But the default metrics were not enough.

We added custom metrics to our Channels so that we could continuously monitor the health of the data flow:

  • Channel Depth: How many messages are waiting in the queue to be processed? If this number rises, it means the consumer is slowing down.
  • Ingest vs. Drain Rate: The rate at which data enters the Channel compared with the rate at which it leaves.
  • Active SSE Connections: The number of online users connected to each node.

These dashboards allowed us to identify bottlenecks before users noticed any slowdown.

The monitoring results were very interesting.

Despite receiving approximately 1,900 messages per second, the architecture was fast enough that our charts showed a maximum of only 15 messages remaining inside a Channel at any given moment.

They were drained almost immediately.

This meant our internal processing latency was effectively near zero.

This shows just how well Channels worked for us. ♥️

Benchmark and Performance

You might be wondering what the result of all this attention to choosing the right technologies — Redis In-Memory, SSE, and Channels — actually was.

The production results exceeded our expectations.

The architecture turned out to be extremely lightweight.

We were able to handle a large amount of traffic with very limited resources:

  • Architecture: A Kubernetes cluster with 6 active Pods.
  • RAM usage: Each Pod consumed only around 55 MB of RAM on average.
  • CPU usage: During peak market traffic, CPU usage per Pod was around 20%.
  • Capacity: This small cluster handled tens of thousands of simultaneous open connections without clients experiencing any noticeable slowdown.

This means that with relatively lightweight infrastructure, we are serving tens of thousands of users without noticeable degradation in quality or latency.

The Tricks That Matter

During implementation, we ran into a few solid walls.

The lessons from those problems can be valuable for any team building a similar system.

The Proxy Trap

When we first deployed the service, we noticed that some users were receiving data with delays, while others were losing the connection entirely.

The problem was not in our code.

It came from the network infrastructure and proxies such as Nginx or enterprise proxies.

SSE uses a long-lived open connection.

Many proxies buffer responses by default before sending them onward, which can break real-time SSE behavior.

Solution: We set:

X-Accel-Buffering: no

We also made sure our proxy configuration allowed long-running streams to pass through without buffering.

The Memory Leak Challenge

The most dangerous scenario in this architecture is when users lose their internet connection or close the browser tab while the server continues writing into their dedicated Channel.

If those Channels are not cleaned up, the server can eventually run into serious memory problems.

Solution: We took the RequestAborted mechanism in HttpContext very seriously.

As soon as the server detects that the client connection has been aborted, we dispose of the corresponding Channel and remove it from the Coordinator.

Scalability

What happens on a busy market day if our load suddenly doubles?

Did we actually design for the worst case?

Solution: We used Kubernetes HPA — Horizontal Pod Autoscaler.

Our system is configured so that as the number of users increases and CPU or RAM consumption rises, Kubernetes automatically increases the number of Pods.

This helps ensure that the system remains ready to serve Mofid users during traffic spikes.

Connection Drops

When the client sends an SSE request, there may be periods when no new data is available.

For example, a stock or fund may temporarily stop trading and therefore have no new price updates.

Long idle connections can also be closed by proxies, load balancers, gateways, or other network infrastructure.

Solution: We took inspiration from heartbeat mechanisms used in network protocols.

Approximately every 20 seconds, the server generates an HB (heartbeat) message.

This helps us verify that the connection is still healthy and also prevents the connection from remaining completely idle.

Conclusion

This project reminded us that in real-time systems, the most complicated tool is not always the best answer.

We kept Kafka where durability and replay mattered.

We chose Redis Pub/Sub for fast and ephemeral fan-out.

In the backend, we used System.Threading.Channels to build a lightweight and low-overhead path for distributing messages.

And on the client side, SSE allowed us to push prices with low latency without requiring a bidirectional protocol.

The key point was understanding the nature of the data before choosing the tools.

A real-time price is replaceable data, not an event that can never be lost.

That decision allowed us to consciously accept limited data loss in places such as Redis Pub/Sub and DropOldest while maintaining freshness, simplicity, and system stability.

In the end, the success of this architecture was not just the result of choosing the right technologies.

It came from the right combination of monitoring, gradual rollout, snapshot/re-sync, heartbeat, feature flags, and close collaboration between the product, backend, frontend, and SRE teams.