JD2022-TU1/main/tools/framework/JD.ElasticSearch/ElasticSearchReporter.cs

137 lines
3.8 KiB
C#

using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Threading;
namespace JD.ElasticSearch
{
public class ElasticSearchReporter: IDisposable
{
#region Constants
private static readonly TimeSpan SleepInterval = TimeSpan.FromSeconds(2);
#endregion
#region Variables
private readonly Thread _processingThread;
private readonly ConcurrentQueue<ElasticSearchLog> _queue = new ConcurrentQueue<ElasticSearchLog>();
private readonly ManualResetEvent _disposedEvent = new ManualResetEvent(false);
private readonly ElasticSearchClient elasticClient;
private DateTime? _shutDownRequestTime;
private bool m_isDisposed;
#endregion
#region Constructor
public ElasticSearchReporter(string [] connectionNodes, string userName = "", string password = "", string defaultIndex = "*")
{
// Setup processing thread
_processingThread = new Thread(RunThread)
{
Name = "ElasticSearch Reporter"
};
_processingThread.Start();
elasticClient = new ElasticSearchClient(connectionNodes, userName, password, defaultIndex);
}
#endregion
#region Methods
public void AddStatistic(ElasticSearchLog log)
{
_queue.Enqueue(log);
}
private void RunThread()
{
while (Wait(SleepInterval))
{
ConsumePendingStatistics();
}
ConsumePendingStatistics();
}
protected virtual bool Wait(TimeSpan interval)
{
return !_disposedEvent.WaitOne(interval, false);
}
public void ConsumePendingStatistics()
{
var statistics = new List<ElasticSearchLog>();
ElasticSearchLog statistic;
const int maxStatsCount = 1000;
while (_queue.TryDequeue(out statistic) && statistics.Count < maxStatsCount)
{
statistics.Add(statistic);
}
if (statistics.Count > 0)
{
SendStatistics(statistics);
}
}
private void SendStatistics(IEnumerable<ElasticSearchLog> statistics)
{
foreach (ElasticSearchLog stat in statistics)
{
if (StatsHaveReachedTimeout)
{
return;
}
var values = stat.Values;
elasticClient.PutDocument(stat.Index, stat.Type, stat.Id, values);
}
}
protected virtual void Dispose(bool disposing)
{
if (disposing)
{
_shutDownRequestTime = DateTime.Now;
_disposedEvent.Set();
if (!_processingThread.Join(TimeSpan.FromSeconds(30)))
{
_processingThread.Abort();
}
}
}
public void Dispose()
{
if (!m_isDisposed)
{
m_isDisposed = true;
Dispose(true);
}
}
#endregion
/// <summary>
/// When closing the application, we are waiting on all the stats to be sent to the server
/// (the thread is not a background thread)
/// But we dont want the application to hang for more than 5 seconds waiting for this to happen.
/// </summary>
private bool StatsHaveReachedTimeout
{
get
{
const int timeAllowedToFlushQueueOnQuit = 5;
return _shutDownRequestTime != null && DateTime.Now.Subtract(_shutDownRequestTime.Value).TotalSeconds > timeAllowedToFlushQueueOnQuit;
}
}
}
}