Skip to main content

Channel, ChannelReader and ChannelWriter to manage data streams in multi-threading environment

I came across Channel class while working with SignalR which looks really interesting. By looking into NuGet packages (https://www.nuget.org/packages/System.Threading.Channels), it seems just 4 months old.

The Channel class provides infrastructure to have multiple reads and write simuletensely through it's Reader and Writer properties.

This is where it is handy in case of SignalR where data streaming needs to be done but is not just limited to that but wherever something needs to be read/write/combination of both in a multi-threading environment.

In my case with SignalR, I had to stream stock data at a regular interval of time.

 public ChannelReader<StockData> StreamStock()  
 {  
   var channel = Channel.CreateUnbounded<StockData>();  
   _stockManager.OnStockData = stockData =>  
   {  
     channel.Writer.TryWrite(stockData);  
   };  
   return channel.Reader;  
 }  

The SignalR keeps return type of ChannelReader<StockData> open so that whatever written in Channel it would be transmitted data through ChannelReader.
The Channel.CreateUnbounded<StockData> creates an instance it also has overridden function to pass allowed number of items which can be stored in the channel which helps out to avoid a lot of buffer usage.

The Reader property of Channel is returned which does not die in SignalR, and whenever anything is written in Channel it is sent to SignalR consumers.

For Writer I am using delegate definition from OnStockData property, whenever the delegate is called from producer class I am calling channel.Writer.TryWrite to pass newly available data.

I have done a little experiment on a console application. Give it a try and check it out. Happy Coding!

 public static class Example  
 {  
   public static async Task<ChannelReader<string>> RunExample()  
   {  
     const int maxMessagesToBuffer = 2;  
     Channel<string> channel = Channel.CreateBounded<string>(maxMessagesToBuffer);  
     var reader = channel.Reader;  
     var writer = channel.Writer;  
     var worker1 = Task.Run(() => ChannelReader(reader,"R1"));  
     var worker2 = Task.Run(() => ChannelReader(reader, "R2"));  
     var worker3 = Task.Run(() => ChannelReader(reader, "R3"));  
     var worker4 = Task.Run(() => ChannelReader(reader, "R4"));  
     await Task.WhenAll(new List<Task> {  
       Task.Run(() => WriteValue(writer, "W1")),  
       Task.Run(() => WriteValue(writer, "W2")),  
       Task.Run(() => WriteValue(writer, "W3")),  
       Task.Run(() => WriteValue(writer, "W4")),  
       Task.Run(() => WriteValue(writer, "W5")),  
     });  
     return reader;  
   }  
   static int index = 0;  
   static SemaphoreSlim @lock = new SemaphoreSlim(1, 1);  
   private async static Task WriteValue(ChannelWriter<string> writer, string writerId)  
   {  
     while (true)  
     {  
       try  
       {  
         await @lock.WaitAsync();  
         if (index < 50)  
         {  
           index++;  
           await writer.WriteAsync($"(index : {index}, Writer: {writerId})");  
         }  
         else  
         {  
           Console.WriteLine($"Writer completed: {writerId}");  
           break;  
         }  
       }  
       finally  
       {  
         @lock.Release();  
       }  
     }  
   }  
   private static async Task ChannelReader(ChannelReader<string> reader, string readerId)  
   {  
     while (await reader.WaitToReadAsync())  
     {  
       while (reader.TryRead(out var message))  
       {  
         Console.WriteLine($"Listener: message : {message}, index dirty read : {index}, Reader: {readerId}");  
         await Task.Delay(250);  
       }  
     }  
   }  
 }  

 public static async Task Main(string[] args)  
 {  
   await Example.RunExample();  
 }  







Comments

Popular posts from this blog

Handling JSON DateTime format on Asp.Net Core

This is a very simple trick to handle JSON date format on AspNet Core by global settings. This can be applicable for the older version as well. In a newer version by default, .Net depends upon Newtonsoft to process any JSON data. Newtonsoft depends upon Newtonsoft.Json.Converters.IsoDateTimeConverter class for processing date which in turns adds timezone for JSON data format. There is a global setting available for same that can be adjusted according to requirement. So, for example, we want to set default formatting to US format, we just need this code. services.AddMvc() .AddJsonOptions(options => { options.SerializerSettings.DateTimeZoneHandling = "MM/dd/yyyy HH:mm:ss"; });

Making FluentValidation compatible with Swagger including Enum or fixed List support

FluentValidation is not directly compatible with Swagger API to validate models. But they do provide an interface through which we can compose Swagger validation manually. That means we look under FluentValidation validators and compose Swagger validator properties to make it compatible. More of all mapping by reading information from FluentValidation and setting it to Swagger Model Schema. These can be done on any custom validation from FluentValidation too just that proper schema property has to be available from Swagger. Custom validation from Enum/List values on FluentValidation using FluentValidation.Validators; using System.Collections.Generic; using System.Linq; using static System.String; /// <summary> /// Validator as per list of items. /// </summary> /// <seealso cref="PropertyValidator" /> public class FixedListValidator : PropertyValidator { /// <summary> /// Gets the valid items /// <...

main method return value

Mainly we used to write "static void main" for entry point in console application. Placement of void denotes return type. In main function we could have "int" too but what does it really mean. "int main" signifies return type as integer. The return type of main function tells about execution status of application. Even if we have specified void as return type then it would be marked as successful program execution. If we mark int as return type then we are able to control the execution status. Now, what is the benefit of making main function as int. Windows OS saves result in  %ERRORLEVEL% environment variable of OS. If we create batch file and execute application through it then we will able to get status and based on result we can trigger something else through batch file. Let's suppose we have created program called TEST.EXE. Batch file script: @echo off REM Execute main program REM TEST.EXE @if  "%ERRORLEVEL%" == "0...

A wrapper implementation for Kendo Grid usage

A wrapper implementation for any heavily used item is always a good practice. Whatever is not written by us and used at a lot of places should be wrapped within specific functionality to keep it future proof and easily changeable. This also encourages DRY principle to keep our common setting at a central place. Kendo UI items are enormous in configuration, one of an issue I find people keep repeating codes for Kendo Grid configuration. They have built very flexible system to have any configuration, but in most of the cases, we do not need all of those complicated configuration. We would try to see a simpler configuration of same. The actual core implementation is bit complex, but we do not have to bother about it once done since the focus is just on usage only. I recommend doing this practice for as simple as jQuery events, form handling or as simple as any notification system. This just won't make things simple but makes codes much more manageable, easy understand, read or open f...

Kendo MVC Grid DataSourceRequest with AutoMapper

Kendo Grid does not work directly with AutoMapper but could be managed by simple trick using mapping through ToDataSourceResult. The solution works fine until different filters are applied. The problems occurs because passed filters refer to view model properties where as database model properties are required after AutoMapper is implemented. So, the plan is to intercept DataSourceRequest  and modify names based on database model. To do that we are going to create implementation of  CustomModelBinderAttribute to catch calls and have our own implementation of DataSourceRequestAttribute from Kendo MVC. I will be using same source code from Kendo but will replace column names for different criteria for sort, filters, group etc. Let's first look into how that will be implemented. public ActionResult GetRoles([MyDataSourceRequest(GridId.RolesUserGrid)] DataSourceRequest request) { if (request == null) { throw new Argume...

Voice control Sony Bravia Television through Alexa

This is my second useful thing done through Alexa after simple implementation of switching on/off light. This is not just applicable to Sony Bravia TVs but any device which can be controlled through HTTP/JSON request or via any other protocol. Hardware prerequisites for making whole thing work are as follows: 1. Sony Bravia Android TV or other devices which can accept input through HTTP or different protocol. 2. Raspberry Pi to keep running program/service. 3. Alexa device. Software prerequisites: 1. Alexa Skill: https://developer.amazon.com/edw/home.html#/skills 2. Lambda: https://console.aws.amazon.com/lambda/home?region=us-east-1#/functions 3. AWS IoT: https://console.aws.amazon.com/iot/home?region=us-east-1 How the whole process would work? Alexa would accept voice commands and converts it to intend to make a request to Lambda function. Lambda function would use converted user-friendly commands to MQTT request on AWS IoT service which would be listened through MQTT ...

Trim text in MVC Core through Model Binder

Trimming text can be done on client side codes, but I believe it is most suitable on MVC Model Binder since it would be at one place on infrastructure level which would be free from any manual intervention of developer. This would allow every post request to be processed and converted to a trimmed string. Let us start by creating Model binder using Microsoft.AspNetCore.Mvc.ModelBinding; using System; using System.Threading.Tasks; public class TrimmingModelBinder : IModelBinder { private readonly IModelBinder FallbackBinder; public TrimmingModelBinder(IModelBinder fallbackBinder) { FallbackBinder = fallbackBinder ?? throw new ArgumentNullException(nameof(fallbackBinder)); } public Task BindModelAsync(ModelBindingContext bindingContext) { if (bindingContext == null) { throw new ArgumentNullException(nameof(bindingContext)); } var valueProviderResult = bindingContext.ValueProvider.GetValue(bin...

LDAP with ASP.Net Identity Core in MVC with project.json

Lightweight Directory Access Protocol (LDAP), the name itself explain it. An application protocol used over an IP network to access the distributed directory information service. The first and foremost thing is to add references for consuming LDAP. This has to be done by adding reference from Global Assembly Cache (GAC) into project.json "frameworks": { "net461": { "frameworkAssemblies": { "System.DirectoryServices": "4.0.0.0", "System.DirectoryServices.AccountManagement": "4.0.0.0" } } }, These  System.DirectoryServices  and  System.DirectoryServices.AccountManagement  references are used to consume LDAP functionality. It is always better to have an abstraction for irrelevant items in consuming part. For an example, the application does not need to know about PrincipalContext or any other dependent items from those two references to make it extensible. So, we can begin wi...

Global exception handling and custom logging in AspNet Core with MongoDB

In this, we would be looking into logging and global exception handling in the AspNet Core application with proper registration of logger and global exception handling. Custom logging The first step is to create a data model that we want to save into DB. Error log Data model These are few properties to do logging which could be extended or reduced based on need. public class ErrorLog { /// <summary> /// Gets or sets the Error log identifier. /// </summary> /// <value> /// The Error log identifier. /// </value> [BsonRepresentation(BsonType.ObjectId)] public ObjectId Id { get; set; /// <summary> /// Gets or sets the date. /// </summary> /// <value> /// The date. /// </value> public DateTime Date { get; set; } /// <summary> /// Gets or sets the thread. /// </summary> ...

Dependency Injection through XML configuration and XML transformation through SlowCheetah

Dependency Injection (DI) is a design pattern to change definition by substituting object without changing code for the application. The most popular DI type is to construct classes based on certain interface and pass actual object on constructor level. What are we trying to do? We will be looking into a way to achieve dependency injection through XML configuration based on  build selection . The implemented classes derived through interface will get switched based on build selection. Where to use it? I really hate making dependencies with something specific which can be changed later on. In my case, Azure environment. I believe Azure is more like a platform where we can host the application rather then integrating the application with Azure. What if client decided to switch to other hosting environment, in that case we got to change every piece of code wherever Azure SDKs are referred. The above one is merely an example. We can use this approach on many other item as wel...