Skip to content

Commit

Permalink
Paralellize broadcast
Browse files Browse the repository at this point in the history
  • Loading branch information
longfin committed Jun 5, 2021
1 parent 3f7bea5 commit 2f4e171
Showing 1 changed file with 26 additions and 23 deletions.
49 changes: 26 additions & 23 deletions Libplanet/Net/Transports/NetMQTransport.cs
Original file line number Diff line number Diff line change
Expand Up @@ -620,34 +620,37 @@ private void DoBroadcast(object sender, NetMQQueueEventArgs<(Address?, Message)>
DateTimeOffset.UtcNow,
_appProtocolVersion);

foreach (BoundPeer peer in peers)
peers.AsParallel().ForAll(peer =>
{
string endpoint = peer.ToNetMQAddress();
if (!_dealers.TryGetValue(peer.Address, out DealerSocket dealer) ||
dealer.IsDisposed)
Task.Run(() =>
{
dealer = new DealerSocket(endpoint);
_dealers[peer.Address] = dealer;
}
else if (dealer.Options.LastEndpoint != endpoint)
{
dealer.Dispose();
dealer = new DealerSocket(endpoint);
_dealers[peer.Address] = dealer;
}
if (!_dealers.TryGetValue(peer.Address, out DealerSocket dealer) ||
dealer.IsDisposed)
{
dealer = new DealerSocket(endpoint);
_dealers[peer.Address] = dealer;
}
else if (dealer.Options.LastEndpoint != endpoint)
{
dealer.Dispose();
dealer = new DealerSocket(endpoint);
_dealers[peer.Address] = dealer;
}

if (!dealer.TrySendMultipartMessage(TimeSpan.FromSeconds(3), message))
{
_logger.Warning(
"Broadcasting timed out. [Peer: {Peer}, Message: {Message}]",
peer,
msg
);
if (!dealer.TrySendMultipartMessage(TimeSpan.FromSeconds(3), message))
{
_logger.Warning(
"Broadcasting timed out. [Peer: {Peer}, Message: {Message}]",
peer,
msg
);

dealer.Dispose();
_dealers.TryRemove(peer.Address, out _);
}
}
dealer.Dispose();
_dealers.TryRemove(peer.Address, out _);
}
});
});
}
catch (Exception exc)
{
Expand Down

0 comments on commit 2f4e171

Please sign in to comment.