Azure / Azure/DotNetty

Memory Leak in Linux CentOS 7.1

Open
#232 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
C#
Stars
4.3k
Forks
1k
PR merge metrics
No merged PRs in 30d

Description

I use dotnetty 0.4.4 in my scheduler project as the network layer,
when i deploy it to production(centos 7.1;8 core cpu;16 G RAM;dotnet 1.1.1)
i found it menory grown up slowing ;by 100-300 bytes per 2-5 seconds, it can grown up to 800-900 MB both client and server side; but the application didn't crash,this is very odd.

client side code:
private void ConnectMaintain()
{
Task.Factory.StartNew( async() => {
//do
//{
if (!IsConnected)
{
try
{
bootstrapChannel = await bootstrap.ConnectAsync(this.serverIP,this.localIP);

}
catch(Exception ex)
{
await this.handler.Disconnect();
OnConnectionFailed?.Invoke(ex);

}
}
await Task.Delay(3000);
//} while (tokenSource.IsCancellationRequested == false && config.AutoRecoverConnect);
}, TaskCreationOptions.LongRunning);
}
public void ConnectToServer()
{
SetConsoleLogger();
group = new MultithreadEventLoopGroup();

bootstrap = new Bootstrap();

bootstrap.Group(group)
.Channel()

.Option(ChannelOption.TcpNodelay, true)
.Option(ChannelOption.SoKeepalive, false)
.Option(ChannelOption.ConnectTimeout, TimeSpan.FromSeconds(2))
.Option(ChannelOption.SoReuseaddr, true)
.Option(ChannelOption.SoLinger, 0)
.Option(ChannelOption.SoBacklog, 200)
.Option(ChannelOption.Allocator, PooledByteBufferAllocator.Default)
.RemoteAddress(this.serverIP)
.LocalAddress(localIP)
.Handler(new ActionChannelInitializer(channel =>
{
IChannelPipeline pipeline = channel.Pipeline;
pipeline.AddLast(new ProtobufVarint32FrameDecoder());
pipeline.AddLast(new ProtobufDecoder(Packet.Parser));
pipeline.AddLast(new ProtobufVarint32LengthFieldPrepender());
pipeline.AddLast(new ProtobufEncoder());
pipeline.AddLast(handler);
}));

this.ConnectMaintain();
}

client side handler code:
internal class JobClientHandler : SimpleChannelInboundHandler
{
IChannelHandlerContext ctx;
internal Action OnJobReceived;
internal Action, int> OnJobsReceived;
public JobClientConfig Config { get; set; }
public int UnixNow
{
get
{
return (int)DateTimeOffset.UtcNow.ToUnixTimeSeconds();
}
}
public override bool IsSharable => true;

public override void ChannelActive(IChannelHandlerContext ctx)
{
this.ctx = ctx;
this.SendPing();
}
public override void ChannelInactive(IChannelHandlerContext context)
{

base.ChannelInactive(context);
}
public override void ExceptionCaught(IChannelHandlerContext context, Exception exception)
{

base.ExceptionCaught(context, exception);
}

public void SendJobState(JobInfo info)
{

ctx.WriteAndFlushAsync(JobPacketUtil.BuildJobPack(info));
}
public async Task Disconnect()
{
if (this.ctx != null)
{
await this.ctx.CloseAsync();
await this.DisconnectAsync(this.ctx);
}
}
protected override void ChannelRead0(IChannelHandlerContext ctx, Packet packet)
{
var msg = packet.Payload;
if (msg.Request.TypeUrl.EndsWith(nameof(PingRequest)))
{
var response = msg.Request.Unpack();
Console.WriteLine($"Receive server response pong:" + response.Timestamp);
var address = ((IPEndPoint)ctx.Channel.LocalAddress);
JobClientInfo clientInfo = null;
Task.Run(async () => {
clientInfo = await ClientUtil.CollectClientInfo(new JobClientConfig { LocalIP = address.Address.MapToIPv4().ToString(), LocalPort = address.Port.ToString(), AutoRecoverConnect = Config.AutoRecoverConnect });
}).Wait();
ctx.WriteAndFlushAsync(JobPacketUtil.BuildClientInfo(clientInfo));

}
else if (msg.Request.TypeUrl.EndsWith(nameof(JobProto)))
{
var job = msg.Request.Unpack();
var info = JobPacketUtil.BuildJobInfo(job);
OnJobReceived?.Invoke(ctx, info);

}
else if (msg.Request.TypeUrl.EndsWith(nameof(JobGroupProto)))
{
var jobGroup = msg.Request.Unpack();
OnJobsReceived?.Invoke(ctx, JobPacketUtil.BuildJobInfos(jobGroup), jobGroup.Policy);

}
else
{
Console.WriteLine("Unknown Packet Message Payload:{0}", msg.Request.TypeUrl);
}
}

void SendPing()
{
var pingRequest = new PingRequest() { Timestamp = new Google.Protobuf.WellKnownTypes.Timestamp() { Seconds = UnixNow } };
var packet = new Packet();
packet.Payload = new Message() { Request = Any.Pack(pingRequest) };
this.ctx.WriteAndFlushAsync(packet);
}
===================================================
server side code:

public class SchedulerServer
{
IChannel bootstrapChannel;
MultithreadEventLoopGroup bossGroup = new MultithreadEventLoopGroup(1);
MultithreadEventLoopGroup workerGroup = new MultithreadEventLoopGroup();
ServerBootstrap bootstrap = new ServerBootstrap();
SchedulerServerHandler handler = new SchedulerServerHandler();
IMailSender mailSender;
readonly SchedulerServerConfig config;

public SchedulerServer(SchedulerServerConfig config,IMailSender mailSender)
{
this.config = config;
this.mailSender = mailSender;
handler.HeartBeatSeconds = config.HeartBeatSeconds<=0?10:config.HeartBeatSeconds;
handler.ZeroClientConnectedAlarm = OnZeroClientConnectedAlarm;
handler.UnKnownClientIDAlarm = OnUnKownClientAlarm;

}
public virtual void OnZeroClientConnectedAlarm(JobInfo job)
{
var mail =string.IsNullOrEmpty( job.AlarmEmail) ? this.config.AlarmEmail:job.AlarmEmail;
if (!string.IsNullOrEmpty(mail))
{
//todo:
}
}
public virtual void OnUnKownClientAlarm(JobInfo job)
{
var mail = string.IsNullOrEmpty(job.AlarmEmail) ? this.config.AlarmEmail:job.AlarmEmail;
if (!string.IsNullOrEmpty(mail))
{
//todo:
}
}
public void StartServer()
{
Process proc = Process.GetCurrentProcess();
proc.EnableRaisingEvents = true;
proc.Exited += Proc_Exited;
InternalLoggerFactory.DefaultFactory.AddProvider(new ConsoleLoggerProvider((s, level) => true, false));
try
{
bootstrap.Group(bossGroup, workerGroup)
.Channel()

.Option(ChannelOption.ConnectTimeout,TimeSpan.FromSeconds(2))
.Option(ChannelOption.SoTimeout, 1000)
.Option(ChannelOption.TcpNodelay,true)
.Option(ChannelOption.SoLinger,0)
.Option(ChannelOption.SoReuseaddr,true)

.Handler(new LoggingHandler("LSTN"))
.Option(ChannelOption.Allocator,PooledByteBufferAllocator.Default)
.ChildHandler(new ActionChannelInitializer(
channel =>
{
IChannelPipeline pipeline = channel.Pipeline;

pipeline.AddLast(new LoggingHandler("CONN"));
pipeline.AddLast(new ProtobufVarint32FrameDecoder());
pipeline.AddLast(new ProtobufDecoder(Packet.Parser));
pipeline.AddLast(new ProtobufVarint32LengthFieldPrepender());
pipeline.AddLast(new ProtobufEncoder());
pipeline.AddLast(handler);
}
))
.ChildOption(ChannelOption.Allocator, PooledByteBufferAllocator.Default);
IPEndPoint endpoint = new IPEndPoint(IPAddress.Parse(config.IPAddress), int.Parse(config.Port));

bootstrapChannel = bootstrap.BindAsync(endpoint).ConfigureAwait(false).GetAwaiter().GetResult();


}
catch
{
throw;
}
}

private void Proc_Exited(object sender, EventArgs e)
{
this.StopServer().Wait();
}

public Task SendJobGroupToClientAsync(DispatchGroupBag bag)
{
return handler.SendJobGroupToClientAsync(bag);
}

public async Task StopServer()
{
try
{

await bootstrapChannel.CloseAsync();
handler.ClearAllChannel();

}
finally {
await Task.WhenAll(bossGroup.ShutdownGracefullyAsync(), workerGroup.ShutdownGracefullyAsync());
}
}
public Task> GetAllClients()
{
return Task.FromResult( handler.GetClients());
}

public Task DisconnectClient(string clientId)
{
return handler.DisconnectClientAsync(clientId);
}
server side handler code:
internal class SchedulerServerHandler : SimpleChannelInboundHandler
{
volatile IChannelGroup ClientChannelGroup;

private Random ran = new Random();
private ConcurrentDictionary clientDic = new ConcurrentDictionary();
internal Action ZeroClientConnectedAlarm;
internal Action UnKnownClientIDAlarm;

public int HeartBeatSeconds { get; internal set; }
public override bool IsSharable
{
get
{
return true;
}
}
public override void ChannelActive(IChannelHandlerContext context)
{
var g = ClientChannelGroup;
if (g == null)
{
lock (this)
{
g = ClientChannelGroup = new DefaultChannelGroup(context.Executor);
HeartBeat(context);
}
}
g.Add(context.Channel);
var clientInfo = new JobClientInfo
{
Id = context.Channel.Id.ToString(),
IPAndPort = context.Channel.RemoteAddressToIPV4()
};
clientDic.TryAdd(context.Channel.Id.ToString(), clientInfo);

}
private void HeartBeat(IChannelHandlerContext context)
{
context.Executor.ScheduleAsync( async (ctx) => {
var c = (ctx as IChannelHandlerContext);

if (!c.Executor.IsShuttingDown && !c.Executor.IsShutdown && !c.Executor.IsTerminated)
{
if (ClientChannelGroup != null)
{
try
{
await ClientChannelGroup.WriteAndFlushAsync(JobPackUtil.BuildPingPack());
}
catch (Exception ex)
{
Console.WriteLine("Exception occur when HeartBeat:"+ex.ToString());
}
}
HeartBeat(c);
}

}, context, TimeSpan.FromSeconds(HeartBeatSeconds));
}
public override void ExceptionCaught(IChannelHandlerContext context, Exception exception)
{
Console.WriteLine($"ExceptionCaught:Channel.Active:{context.Channel.Active},Channel.IsWritable:{context.Channel.IsWritable},Channel.IsWritable:{context.Channel.ToString()}");
if (context.Channel.Active == false)
{
context.Channel.DeregisterAsync();
context.Channel.CloseAsync().Wait();
ClientChannelGroup.Remove(context.Channel);
}
}
private IChannel RandomChannel()
{
var channels = ClientChannelGroup.ToList();
var randomIndex = ran.Next(0, channels.Count);
return channels[randomIndex];
}
private async Task SendToRandomClientAsync(Packet packet)
{
var channel = RandomChannel();
await channel.WriteAndFlushAsync(packet);
}
public async Task SendJobToClientAsync(JobInfo info)
{
if (ClientChannelGroup == null)
{
ZeroClientConnectedAlarm?.Invoke(info);
return;
}
Packet packet = JobPackUtil.BuildJobPack(info);
if (string.IsNullOrEmpty(info.ClientIPAndPort) == false)
{
if (ClientChannelGroup.Any(i => i.RemoteAddressToIPV4() == info.ClientIPAndPort) == false)
{
UnKnownClientIDAlarm?.Invoke(info);
return;
}
await ClientChannelGroup.WriteAndFlushAsync(packet, new ClientIPMatcher(info.ClientIPAndPort));
}
else
{
await SendToRandomClientAsync(packet);
}
}

public async Task SendJobGroupToClientAsync(DispatchGroupBag bag)
{
if (ClientChannelGroup == null || ClientChannelGroup.Any() == false)
{
ZeroClientConnectedAlarm?.Invoke(bag.Items.First());
return;
}
//指定客戶端發送
if (bag.Items.Any(i => string.IsNullOrEmpty(i.ClientIPAndPort) == false))
{
var jobInfo = bag.Items.First();
if (ClientChannelGroup.Any(i => i.RemoteAddressToIPV4() == jobInfo.ClientIPAndPort) == false)
{
UnKnownClientIDAlarm?.Invoke(jobInfo);
return;
}
Packet packet = JobPackUtil.BuildJobGroupPack(bag.Items, (int)bag.Policy);
await ClientChannelGroup.WriteAndFlushAsync(packet, new ClientIPMatcher(bag.Items.First().ClientIPAndPort));
return;
}
if (bag.Policy == DispatchPolicy.Broadcast)
{
Packet packet = JobPackUtil.BuildJobGroupPack(bag.Items, (int)bag.Policy);
await ClientChannelGroup.WriteAndFlushAsync(packet);
return;
}
List tasks = new List();
if (bag.Policy == DispatchPolicy.RandomBoth)
{
foreach (var job in bag.Items)
{
tasks.Add( SendJobToClientAsync(job));
}
}

if (bag.Policy == DispatchPolicy.OneClientAsync ||
bag.Policy == DispatchPolicy.OneClientSequence ||
bag.Policy == DispatchPolicy.OneClientSequenceBreakOnFails
)
{
var channel = RandomChannel();
Packet packet = JobPackUtil.BuildJobGroupPack(bag.Items,(int)bag.Policy);
tasks.Add( channel.WriteAndFlushAsync(packet));
}

if (bag.Policy == DispatchPolicy.OneByOne)
{
var channels = ClientChannelGroup.ToList();
var jobs = new Stack(bag.Items.ToList());
while (jobs.Any())
{
foreach (var c in channels)
{
if (jobs.Any() == false)
{
break;
}
var job = jobs.Pop();
Packet packet = JobPackUtil.BuildJobPack(job);
tasks.Add( c.WriteAndFlushAsync(packet));
}
}
}
if (tasks.Any())
{
await Task.WhenAll(tasks.ToArray());
}

}

public async Task DisconnectClientAsync(string clientId)
{
if (ClientChannelGroup != null)
{
var clientChannel = ClientChannelGroup.Where(i => i.Id.ToString() == clientId).FirstOrDefault();
if (clientChannel != null)
{
await clientChannel.DisconnectAsync();
}
}
}

public List GetClients()
{
return clientDic.Values.ToList();

}
public void ClearAllChannel()
{
if (ClientChannelGroup != null)
{
ClientChannelGroup.Clear();
}
}
public override void ChannelInactive(IChannelHandlerContext context)
{
ClientChannelGroup.Remove(context.Channel);
JobClientInfo removed = null;
clientDic.TryRemove(context.Channel.Id.ToString(), out removed);
}

protected override void ChannelRead0(IChannelHandlerContext ctx, Packet packet)
{
var msg = packet.Payload;
if (msg.Request.TypeUrl.EndsWith(nameof(PingRequest)))
{
var request = msg.Request.Unpack();
ctx.WriteAndFlushAsync(packet);

}
if (msg.Request.TypeUrl.EndsWith(nameof(ClientProto)))
{
var client = msg.Request.Unpack();
var info = clientDic.Values.Where(i => i.IPAndPort == client.IPAndPort).FirstOrDefault();
if (info != null)
{
info.InstanceMemory = client.InstanceMemory;
info.InstanceThreads = client.InstanceThreads;
info.TotalFreeMemory = client.TotalFreeMemory;
info.TotalMemory = client.TotalMemory;
info.CpuCount = client.CpuCount;
info.CpuUsedPercent = client.CpuUsedPercent;
info.AutoRecoverConnect = client.AutoRecoverConnect;
info.UpTime = client.UpTime;
info.ProcessId = client.ProcessId;

}
}
}

}
}

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.