This commit is contained in:
@ -24,10 +24,8 @@ namespace OutputServiceTSDB.RabbitMQ
|
||||
{
|
||||
var factory = new ConnectionFactory { HostName = "localhost" };
|
||||
|
||||
// create connection
|
||||
_connection = factory.CreateConnection();
|
||||
|
||||
// create channel
|
||||
_channel = _connection.CreateModel();
|
||||
|
||||
_channel.ExchangeDeclare("demo.exchange", ExchangeType.Topic);
|
||||
@ -45,35 +43,21 @@ namespace OutputServiceTSDB.RabbitMQ
|
||||
var consumer = new EventingBasicConsumer(_channel);
|
||||
consumer.Received += (ch, ea) =>
|
||||
{
|
||||
// received message
|
||||
var content = System.Text.Encoding.UTF8.GetString(ea.Body);
|
||||
|
||||
// handle the received message
|
||||
HandleMessage(content);
|
||||
_channel.BasicAck(ea.DeliveryTag, false);
|
||||
};
|
||||
|
||||
consumer.Shutdown += OnConsumerShutdown;
|
||||
consumer.Registered += OnConsumerRegistered;
|
||||
consumer.Unregistered += OnConsumerUnregistered;
|
||||
consumer.ConsumerCancelled += OnConsumerConsumerCancelled;
|
||||
|
||||
_channel.BasicConsume("demo.queue.log", false, consumer);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
private void HandleMessage(string content)
|
||||
{
|
||||
// we just print this message
|
||||
_logger.LogInformation($"consumer received {content}");
|
||||
}
|
||||
|
||||
private void OnConsumerConsumerCancelled(object sender, ConsumerEventArgs e) { }
|
||||
private void OnConsumerUnregistered(object sender, ConsumerEventArgs e) { }
|
||||
private void OnConsumerRegistered(object sender, ConsumerEventArgs e) { }
|
||||
private void OnConsumerShutdown(object sender, ShutdownEventArgs e) { }
|
||||
private void RabbitMQ_ConnectionShutdown(object sender, ShutdownEventArgs e) { }
|
||||
|
||||
public override void Dispose()
|
||||
{
|
||||
_channel.Close();
|
||||
|
Reference in New Issue
Block a user