I am using RX 2.2.5
In the sample below, I expect the handler (the delegate to Subscribe()) to run on a new thread however when running the app all of the 10 numbers are consumed on the same thread one after another.
Console.WriteLine("Main Thread: {0}", Thread.CurrentThread.ManagedThreadId);
var source = Observable.Range(1, 10, TaskPoolScheduler.Default)
.ObserveOn(NewThreadScheduler.Default)
.Subscribe(n =>
{
Thread.Sleep(1000);
Console.WriteLine(
"Value: {0} Thread: {1} IsPool: {2}",
n, Thread.CurrentThread.ManagedThreadId, Thread.CurrentThread.IsThreadPoolThread);
});
Output:
Main Thread: 1
Value: 1 on Thread: 4 IsPool: False
Value: 2 on Thread: 4 IsPool: False
Value: 3 on Thread: 4 IsPool: False
Value: 4 on Thread: 4 IsPool: False
Value: 5 on Thread: 4 IsPool: False
Value: 6 on Thread: 4 IsPool: False
Value: 7 on Thread: 4 IsPool: False
Value: 8 on Thread: 4 IsPool: False
Value: 9 on Thread: 4 IsPool: False
Value: 10 on Thread: 4 IsPool: False
The fact that they run sequentially is a mystery as well since I am using the TaskPoolScheduler to generate the numbers.
Even if I replace the NewThreadScheduler with TaskPoolScheduler or ThreadPoolScheduler I still get one thread and the more interesting part is that in both those cases the Thread.CurrentThread.IsThreadPoolThread is False.
I cannot explain this behaviour as when I look at the ThreadPoolScheduler I see:
public override IDisposable Schedule<TState>(TState state, Func<IScheduler, TState, IDisposable> action)
{
if (action == null)
throw new ArgumentNullException("action");
SingleAssignmentDisposable d = new SingleAssignmentDisposable();
ThreadPool.QueueUserWorkItem((WaitCallback) (_ =>
{
if (d.IsDisposed)
return;
d.Disposable = action((IScheduler) this, state);
}), (object) null);
return (IDisposable) d;
}
I can clearly see ThreadPool.QueueUserWorkItem... so why IsPool == False?
What am I missing here?