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
29 changes: 27 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.Linq;
using System.Net;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using log4net;

Expand All @@ -13,9 +15,32 @@ 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 = new List<Task<string>>();
var timeoutTask = Task.Delay(timeout);

for (var i = 0; i < ReplicaAddresses.Length; ++i)
{
var request = CreateRequest(ReplicaAddresses[i] + "?query=" + query);
tasks.Add(ProcessRequestAsync(request));
}

while (tasks.Count > 0)
{
var completedTask = await Task.WhenAny(tasks.Concat(new[] {timeoutTask}));
if (completedTask == timeoutTask)
throw new TimeoutException();
try
{
return await (Task<string>)completedTask;
}
catch
{
tasks.Remove((Task<string>)completedTask);
}
}
throw new TimeoutException();
}

protected override ILog Log => LogManager.GetLogger(typeof(ParallelClusterClient));
Expand Down
42 changes: 40 additions & 2 deletions homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
Expand All @@ -9,13 +10,50 @@ namespace ClusterClient.Clients
{
public class RoundRobinClusterClient : ClusterClientBase
{
private readonly ReplicasStatistics _replicasStatistics;
public RoundRobinClusterClient(string[] replicaAddresses) : base(replicaAddresses)
{
_replicasStatistics = new ReplicasStatistics(replicaAddresses);
}

public override Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
public override async Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
{
throw new NotImplementedException();
var timer = Stopwatch.StartNew();
var sortedReplicas = ReplicaAddresses
.OrderBy(x => _replicasStatistics.GetStats(x))
.ToArray();
var sentTasks = new Dictionary<Task<string>, (string, DateTime)>();
for (int i = 0; i < sortedReplicas.Length; i++)
{
var remainingTime = timeout - timer.Elapsed;
if (remainingTime <= TimeSpan.Zero)
throw new TimeoutException();
var timeoutForReplica = TimeSpan.FromTicks(remainingTime.Ticks / (sortedReplicas.Length - i));
if (timeoutForReplica <= TimeSpan.Zero)
throw new TimeoutException();
var timeoutTask = Task.Delay(timeoutForReplica);
Comment on lines +28 to +34

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

remainingTime и timeoutForReplica могут оказаться отрицательными. Task.Delay может или выбросить исключение, или, с маленьким шансом, превратиться в никогда не заканчивающуюся таску

var request = CreateRequest(sortedReplicas[i] + "?query=" + query);
var task = ProcessRequestAsync(request);
sentTasks.Add(task, (address: sortedReplicas[i], startTime: DateTime.Now));
var requestResult = await Task.WhenAny(task, timeoutTask);
if (requestResult != timeoutTask)
Comment on lines +38 to +39

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

если requestResult == timeoutTask, то статистика для реплики не обновится

{
var replicaData = sentTasks[task];
var time = DateTime.Now - replicaData.Item2;
var result = (Task<string>)requestResult;
if (result.Status == TaskStatus.RanToCompletion)
{
_replicasStatistics.UpdateStats(replicaData.Item1, time.TotalMilliseconds);
sentTasks.Remove(task);
return result.Result;
}
_replicasStatistics.UpdateStats(replicaData.Item1, timeout.TotalMilliseconds);
sentTasks.Remove(task);
}
else
_replicasStatistics.UpdateStats(sortedReplicas[i], timeoutForReplica.TotalMilliseconds);
}
throw new TimeoutException();
}

protected override ILog Log => LogManager.GetLogger(typeof(RoundRobinClusterClient));
Expand Down
48 changes: 45 additions & 3 deletions homework 2/ClusterClient/Clients/SmartClusterClient.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
Expand All @@ -9,15 +10,56 @@ namespace ClusterClient.Clients
{
public class SmartClusterClient : ClusterClientBase
{
private readonly ReplicasStatistics _replicasStatistics;

public SmartClusterClient(string[] replicaAddresses) : base(replicaAddresses)
{
_replicasStatistics = new ReplicasStatistics(replicaAddresses);
}

public override Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
public override async Task<string> ProcessRequestAsync(string query, TimeSpan timeout)
{
throw new NotImplementedException();
var timer = Stopwatch.StartNew();
var sortedReplicas = ReplicaAddresses
.OrderBy(x => _replicasStatistics.GetStats(x))
.ToArray();
var requests = new Dictionary<Task<string>, (string, DateTime)>();
for (int i = 0; i < sortedReplicas.Length; i++)
{
var remainingTime = timeout - timer.Elapsed;
if (remainingTime <= TimeSpan.Zero)
throw new TimeoutException();
var timeoutForReplica = TimeSpan.FromTicks(remainingTime.Ticks / (sortedReplicas.Length - i));
if (timeoutForReplica <= TimeSpan.Zero)
throw new TimeoutException();
var timeoutTask = Task.Delay(timeoutForReplica);
var request = CreateRequest(sortedReplicas[i] + "?query=" + query);
requests.Add(ProcessRequestAsync(request), (sortedReplicas[i], DateTime.Now));
var requestResult = await Task.WhenAny(requests.Keys.Concat(new[] { timeoutTask }));
if (requestResult != timeoutTask)
{
var result = (Task<string>)requestResult;
var replicaData = requests[result];
var time = DateTime.Now - replicaData.Item2;
if (result.Status == TaskStatus.RanToCompletion)
{
_replicasStatistics.UpdateStats(replicaData.Item1, time.TotalMilliseconds);
requests.Remove(result);
foreach (var badRequest in requests)
_replicasStatistics.UpdateStats(badRequest.Value.Item1, timeoutForReplica.TotalMilliseconds);
requests.Clear();
return result.Result;

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Если все реплики будут тормозить, то соответствующие таски так и останутся в requests, и статистика для них не обновится

}
_replicasStatistics.UpdateStats(replicaData.Item1, time.TotalMilliseconds);
requests.Remove(result);
}
else
_replicasStatistics.UpdateStats(sortedReplicas[i], timeoutForReplica.TotalMilliseconds);
}

throw new TimeoutException();
}

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

namespace ClusterClient;

public class ReplicasStatistics
{
private const double ConfidenceFactor = 0.8;
private readonly ConcurrentDictionary<string, double> _replicaStatistics;

public ReplicasStatistics(string[] replicaAddresses)
{
_replicaStatistics = new ConcurrentDictionary<string, double>();
foreach (var replicaAddress in replicaAddresses)
_replicaStatistics.TryAdd(replicaAddress, 0);
}

public void UpdateStats(string replicaAddress, double workTime)
{
_replicaStatistics.AddOrUpdate(
replicaAddress,
workTime,
(_, previousValue) => previousValue * ConfidenceFactor + workTime * (1 - ConfidenceFactor)
);
}

public double GetStats(string replicaAddress)
{
return _replicaStatistics.TryGetValue(replicaAddress, out double value) ? value : 0;
}
}