Support passive queue declaration in amqp_consumer (#5831)
This commit is contained in:
parent
b5cd9a9ff2
commit
e141518cf0
|
@ -27,7 +27,7 @@ The following defaults are known to work with RabbitMQ:
|
||||||
# username = ""
|
# username = ""
|
||||||
# password = ""
|
# password = ""
|
||||||
|
|
||||||
## Exchange to declare and consume from.
|
## Name of the exchange to declare. If unset, no exchange will be declared.
|
||||||
exchange = "telegraf"
|
exchange = "telegraf"
|
||||||
|
|
||||||
## Exchange type; common types are "direct", "fanout", "topic", "header", "x-consistent-hash".
|
## Exchange type; common types are "direct", "fanout", "topic", "header", "x-consistent-hash".
|
||||||
|
@ -49,7 +49,11 @@ The following defaults are known to work with RabbitMQ:
|
||||||
## AMQP queue durability can be "transient" or "durable".
|
## AMQP queue durability can be "transient" or "durable".
|
||||||
queue_durability = "durable"
|
queue_durability = "durable"
|
||||||
|
|
||||||
## Binding Key
|
## If true, queue will be passively declared.
|
||||||
|
# queue_passive = false
|
||||||
|
|
||||||
|
## A binding between the exchange and queue using this binding key is
|
||||||
|
## created. If unset, no binding is created.
|
||||||
binding_key = "#"
|
binding_key = "#"
|
||||||
|
|
||||||
## Maximum number of messages server should give to the worker.
|
## Maximum number of messages server should give to the worker.
|
||||||
|
|
|
@ -41,6 +41,7 @@ type AMQPConsumer struct {
|
||||||
// Queue Name
|
// Queue Name
|
||||||
Queue string `toml:"queue"`
|
Queue string `toml:"queue"`
|
||||||
QueueDurability string `toml:"queue_durability"`
|
QueueDurability string `toml:"queue_durability"`
|
||||||
|
QueuePassive bool `toml:"queue_passive"`
|
||||||
|
|
||||||
// Binding Key
|
// Binding Key
|
||||||
BindingKey string `toml:"binding_key"`
|
BindingKey string `toml:"binding_key"`
|
||||||
|
@ -101,7 +102,7 @@ func (a *AMQPConsumer) SampleConfig() string {
|
||||||
# username = ""
|
# username = ""
|
||||||
# password = ""
|
# password = ""
|
||||||
|
|
||||||
## Exchange to declare and consume from.
|
## Name of the exchange to declare. If unset, no exchange will be declared.
|
||||||
exchange = "telegraf"
|
exchange = "telegraf"
|
||||||
|
|
||||||
## Exchange type; common types are "direct", "fanout", "topic", "header", "x-consistent-hash".
|
## Exchange type; common types are "direct", "fanout", "topic", "header", "x-consistent-hash".
|
||||||
|
@ -123,7 +124,11 @@ func (a *AMQPConsumer) SampleConfig() string {
|
||||||
## AMQP queue durability can be "transient" or "durable".
|
## AMQP queue durability can be "transient" or "durable".
|
||||||
queue_durability = "durable"
|
queue_durability = "durable"
|
||||||
|
|
||||||
## Binding Key.
|
## If true, queue will be passively declared.
|
||||||
|
# queue_passive = false
|
||||||
|
|
||||||
|
## A binding between the exchange and queue using this binding key is
|
||||||
|
## created. If unset, no binding is created.
|
||||||
binding_key = "#"
|
binding_key = "#"
|
||||||
|
|
||||||
## Maximum number of messages server should give to the worker.
|
## Maximum number of messages server should give to the worker.
|
||||||
|
@ -286,6 +291,7 @@ func (a *AMQPConsumer) connect(amqpConf *amqp.Config) (<-chan amqp.Delivery, err
|
||||||
return nil, fmt.Errorf("Failed to open a channel: %s", err)
|
return nil, fmt.Errorf("Failed to open a channel: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if a.Exchange != "" {
|
||||||
var exchangeDurable = true
|
var exchangeDurable = true
|
||||||
switch a.ExchangeDurability {
|
switch a.ExchangeDurability {
|
||||||
case "transient":
|
case "transient":
|
||||||
|
@ -309,27 +315,18 @@ func (a *AMQPConsumer) connect(amqpConf *amqp.Config) (<-chan amqp.Delivery, err
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
var queueDurable = true
|
|
||||||
switch a.QueueDurability {
|
|
||||||
case "transient":
|
|
||||||
queueDurable = false
|
|
||||||
default:
|
|
||||||
queueDurable = true
|
|
||||||
}
|
}
|
||||||
|
|
||||||
q, err := ch.QueueDeclare(
|
q, err := declareQueue(
|
||||||
a.Queue, // queue
|
ch,
|
||||||
queueDurable, // durable
|
a.Queue,
|
||||||
false, // delete when unused
|
a.QueueDurability,
|
||||||
false, // exclusive
|
a.QueuePassive)
|
||||||
false, // no-wait
|
|
||||||
nil, // arguments
|
|
||||||
)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("Failed to declare a queue: %s", err)
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if a.BindingKey != "" {
|
||||||
err = ch.QueueBind(
|
err = ch.QueueBind(
|
||||||
q.Name, // queue
|
q.Name, // queue
|
||||||
a.BindingKey, // binding-key
|
a.BindingKey, // binding-key
|
||||||
|
@ -340,6 +337,7 @@ func (a *AMQPConsumer) connect(amqpConf *amqp.Config) (<-chan amqp.Delivery, err
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("Failed to bind a queue: %s", err)
|
return nil, fmt.Errorf("Failed to bind a queue: %s", err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
err = ch.Qos(
|
err = ch.Qos(
|
||||||
a.PrefetchCount,
|
a.PrefetchCount,
|
||||||
|
@ -402,6 +400,48 @@ func declareExchange(
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func declareQueue(
|
||||||
|
channel *amqp.Channel,
|
||||||
|
queueName string,
|
||||||
|
queueDurability string,
|
||||||
|
queuePassive bool,
|
||||||
|
) (*amqp.Queue, error) {
|
||||||
|
var queue amqp.Queue
|
||||||
|
var err error
|
||||||
|
|
||||||
|
var queueDurable = true
|
||||||
|
switch queueDurability {
|
||||||
|
case "transient":
|
||||||
|
queueDurable = false
|
||||||
|
default:
|
||||||
|
queueDurable = true
|
||||||
|
}
|
||||||
|
|
||||||
|
if queuePassive {
|
||||||
|
queue, err = channel.QueueDeclarePassive(
|
||||||
|
queueName, // queue
|
||||||
|
queueDurable, // durable
|
||||||
|
false, // delete when unused
|
||||||
|
false, // exclusive
|
||||||
|
false, // no-wait
|
||||||
|
nil, // arguments
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
queue, err = channel.QueueDeclare(
|
||||||
|
queueName, // queue
|
||||||
|
queueDurable, // durable
|
||||||
|
false, // delete when unused
|
||||||
|
false, // exclusive
|
||||||
|
false, // no-wait
|
||||||
|
nil, // arguments
|
||||||
|
)
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("error declaring queue: %v", err)
|
||||||
|
}
|
||||||
|
return &queue, nil
|
||||||
|
}
|
||||||
|
|
||||||
// Read messages from queue and add them to the Accumulator
|
// Read messages from queue and add them to the Accumulator
|
||||||
func (a *AMQPConsumer) process(ctx context.Context, msgs <-chan amqp.Delivery, ac telegraf.Accumulator) {
|
func (a *AMQPConsumer) process(ctx context.Context, msgs <-chan amqp.Delivery, ac telegraf.Accumulator) {
|
||||||
a.deliveries = make(map[telegraf.TrackingID]amqp.Delivery)
|
a.deliveries = make(map[telegraf.TrackingID]amqp.Delivery)
|
||||||
|
|
|
@ -33,7 +33,7 @@ For an introduction to AMQP see:
|
||||||
# exchange_type = "topic"
|
# exchange_type = "topic"
|
||||||
|
|
||||||
## If true, exchange will be passively declared.
|
## If true, exchange will be passively declared.
|
||||||
# exchange_declare_passive = false
|
# exchange_passive = false
|
||||||
|
|
||||||
## Exchange durability can be either "transient" or "durable".
|
## Exchange durability can be either "transient" or "durable".
|
||||||
# exchange_durability = "durable"
|
# exchange_durability = "durable"
|
||||||
|
|
|
@ -92,7 +92,7 @@ var sampleConfig = `
|
||||||
# exchange_type = "topic"
|
# exchange_type = "topic"
|
||||||
|
|
||||||
## If true, exchange will be passively declared.
|
## If true, exchange will be passively declared.
|
||||||
# exchange_declare_passive = false
|
# exchange_passive = false
|
||||||
|
|
||||||
## Exchange durability can be either "transient" or "durable".
|
## Exchange durability can be either "transient" or "durable".
|
||||||
# exchange_durability = "durable"
|
# exchange_durability = "durable"
|
||||||
|
|
|
@ -78,6 +78,10 @@ func Connect(config *ClientConfig) (*client, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *client) DeclareExchange() error {
|
func (c *client) DeclareExchange() error {
|
||||||
|
if c.config.exchange == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
var err error
|
var err error
|
||||||
if c.config.exchangePassive {
|
if c.config.exchangePassive {
|
||||||
err = c.channel.ExchangeDeclarePassive(
|
err = c.channel.ExchangeDeclarePassive(
|
||||||
|
|
Loading…
Reference in New Issue