Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 31 additions & 2 deletions homework 2/ClusterClient/Clients/ParallelClusterClient.cs
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
using System;
using System.Collections.Generic;
using System.ComponentModel.DataAnnotations;
using System.Linq;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using log4net;

Expand All @@ -13,9 +15,36 @@ public ParallelClusterClient(string[] replicaAddresses) : base(replicaAddresses)
{
}

public override Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
public override async Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
{
throw new NotImplementedException();
var tasks = ReplicaAddresses.Select(async replica =>
{
var webRequest = CreateRequest(replica + "?query=" + query);
Log.InfoFormat($"Processing {webRequest.RequestUri}");
return await ProcessRequestAsync(webRequest);
}).ToList();

var timeoutTask = Task.Delay(timeout);

while (tasks.Count != 0)
{
var completed = await Task.WhenAny(tasks.Append(timeoutTask));

if (completed == timeoutTask)
throw new TimeoutException();

var result = (Task<string>)completed;
tasks.Remove(result);
try
{
return await result;
}
catch (Exception e)
{

}
Comment on lines +38 to +45

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Тут можно было сделать чуть проще, перед await result просто проврять на result.IsFault, получилось бы чуть короче и без "выбрасывания" исключения.
Но и так тоже можно)

}
throw new Exception();
}

protected override ILog Log => LogManager.GetLogger(typeof(ParallelClusterClient));
Expand Down
34 changes: 34 additions & 0 deletions homework 2/ClusterClient/Clients/ReplicaStatistics.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
using System.Collections.Concurrent;
using System.Linq;

namespace ClusterClient.Clients
{
public static class ReplicaStatistics
{
private static readonly ConcurrentDictionary<string, ConcurrentBag<long>> _responseTimes = new();
private const int MaxStoredTimes = 100;

public static void RecordResponseTime(string replicaAddress, long elapsedMilliseconds)
{
var bag = _responseTimes.GetOrAdd(replicaAddress, _ => new ConcurrentBag<long>());
bag.Add(elapsedMilliseconds);

var newBag = new ConcurrentBag<long>(bag.TakeLast(MaxStoredTimes));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ConcurrentBag не гарантирует порядок обхода, а TakeLast работает на текущей "моментальной" выборке, которая может измениться. Для хранения последних N значений лучше подошёл бы потокобезопасный ограниченный буфер (например, ConcurrentQueue с проверкой длины и TryDequeue при превышении).

_responseTimes.TryUpdate(replicaAddress, newBag, bag);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

TryUpdate не гарантирует, что значение в словаре обновится, другой поток мог быть первым и TryUpdate вернет false.

Для текущей задачи с "приблизительной" статистикой, потерять данные не критично, но в реальных задачах это может стать источником "плавающих багов")

}

private static double GetAverageResponseTime(string replicaAddress)
{
if (!_responseTimes.TryGetValue(replicaAddress, out var times) || times.IsEmpty)
{
return double.MaxValue;
}
return times.Average();
}

public static string[] OrderBySpeed(string[] replicaAddresses)
{
return replicaAddresses.OrderBy(GetAverageResponseTime).ToArray();
}
}
}
36 changes: 30 additions & 6 deletions homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs
Original file line number Diff line number Diff line change
@@ -1,7 +1,5 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Diagnostics;
using System.Threading.Tasks;
using log4net;

Expand All @@ -13,11 +11,37 @@ public RoundRobinClusterClient(string[] replicaAddresses) : base(replicaAddresse
{
}

public override Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
public override async Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
{
throw new NotImplementedException();
var orderedReplicas = ReplicaStatistics.OrderBySpeed(ReplicaAddresses);

var deadline = DateTime.UtcNow + timeout;
var addressesLeft = orderedReplicas.Length;

foreach (var replica in orderedReplicas)
{
var remaining = deadline - DateTime.UtcNow;
var replicaTimeout = remaining / addressesLeft;
Comment on lines +23 to +24

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Тут не хватает проверки: что replicaTimeout не станет отрицательным (к примеру из-за отрицательности remaining).

Если replicaTimeout станет отрицательным и мы попытаемся выполнить Task.Delay(replicaTimeout) то произойдет исключение ArgumentOutOfRangeException


var webRequest = CreateRequest(replica + "?query=" + query);
Log.InfoFormat($"Processing {webRequest.RequestUri}");

var timer = Stopwatch.StartNew();
var resultTask = ProcessRequestAsync(webRequest);

await Task.WhenAny(resultTask, Task.Delay(replicaTimeout));
timer.Stop();

if (resultTask.IsCompletedSuccessfully)
{
ReplicaStatistics.RecordResponseTime(replica, timer.ElapsedMilliseconds);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Здесь и в SmartClusterClient статистика записывается только в случае успешного завершения, хотя лучше записыать ее в любом случае, просто "упавшим" или тормозившим репликам добавлять штраф, чтобы они гарантировано при проблемах уходили вниз списка реплик)

return await resultTask;
}
addressesLeft--;
}
throw new TimeoutException();
}

protected override ILog Log => LogManager.GetLogger(typeof(RoundRobinClusterClient));
}
}
}
74 changes: 68 additions & 6 deletions homework 2/ClusterClient/Clients/SmartClusterClient.cs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
using log4net;

Expand All @@ -12,12 +12,74 @@ public class SmartClusterClient : ClusterClientBase
public SmartClusterClient(string[] replicaAddresses) : base(replicaAddresses)
{
}

public override Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
public override async Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
{
throw new NotImplementedException();
}
var orderedReplicas = ReplicaStatistics.OrderBySpeed(ReplicaAddresses);

var deadline = DateTime.UtcNow + timeout;
var runningTasks = new List<(Task<string> task, string address, Stopwatch timer)>();
var addressesLeft = orderedReplicas.Length;

foreach (var replica in orderedReplicas)
{
var remaining = deadline - DateTime.UtcNow;
var replicaTimeout = remaining / addressesLeft;
Comment on lines +26 to +27

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Здесь так же не хватает проверки, не стал ли replicaTimeout отрицательным или = 0

if (remaining <= TimeSpan.Zero)
                     throw new TimeoutException();


var webRequest = CreateRequest(replica + "?query=" + query);
Log.InfoFormat($"Processing {webRequest.RequestUri}");

var timer = Stopwatch.StartNew();
var task = ProcessRequestAsync(webRequest);
runningTasks.Add((task, replica, timer));

var finished = await WaitOneAsync(runningTasks, replicaTimeout);
if(finished != null)
{
var (resultTask, resultReplica, resultTimer) = finished.Value;
resultTimer.Stop();
ReplicaStatistics.RecordResponseTime(resultReplica, resultTimer.ElapsedMilliseconds);
return resultTask.Result;
}
addressesLeft--;
}

while (runningTasks.Count > 0)
{
var remaining = deadline - DateTime.UtcNow;
if (remaining <= TimeSpan.Zero)
throw new TimeoutException();

var finished = await WaitOneAsync(runningTasks, remaining);
if (finished == null) continue;
var (resultTask, resultReplica, resultTimer) = finished.Value;
resultTimer.Stop();
ReplicaStatistics.RecordResponseTime(resultReplica, resultTimer.ElapsedMilliseconds);
return resultTask.Result;
}
throw new TimeoutException();
}

protected override ILog Log => LogManager.GetLogger(typeof(SmartClusterClient));

private async Task<(Task<string>, string, Stopwatch)?> WaitOneAsync
(List<(Task<string> task, string address, Stopwatch timer)> runningTasks, TimeSpan timeout)
{
if (timeout <= TimeSpan.Zero)
return null;

var delayTask = Task.Delay(timeout);
var tasksList = runningTasks.Select(Task (x) => x.task).Append(delayTask).ToList();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Можно упростить убрав указание типа:
runningTasks.Select(x => x.task).Append(delayTask).ToList()

var completed = await Task.WhenAny(tasksList);

if (completed == delayTask)
return null;

var finishedTask = (Task<string>)completed;
var result = runningTasks.First(x => x.task == finishedTask);
runningTasks.Remove(result);

return finishedTask.IsCompletedSuccessfully ? (finishedTask, result.address, result.timer) : null;
}
}
}
}