fix close channel by queue redecleration (#162)

This commit is contained in:
Yaron Schneider 2019-12-23 16:17:00 -08:00 committed by GitHub
parent 98999a6044
commit 94d9988642
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
1 changed files with 2 additions and 7 deletions

View File

@ -57,12 +57,7 @@ func (r *RabbitMQ) Init(metadata bindings.Metadata) error {
}
func (r *RabbitMQ) Write(req *bindings.WriteRequest) error {
q, err := r.channel.QueueDeclare(r.metadata.QueueName, r.metadata.Durable, r.metadata.DeleteWhenUnused, false, false, nil)
if err != nil {
return err
}
err = r.channel.Publish("", q.Name, false, false, amqp.Publishing{
err := r.channel.Publish("", r.metadata.QueueName, false, false, amqp.Publishing{
ContentType: "text/plain",
Body: req.Data,
})
@ -95,7 +90,7 @@ func (r *RabbitMQ) Read(handler func(*bindings.ReadResponse) error) error {
msgs, err := r.channel.Consume(
q.Name,
"",
true,
false,
false,
false,
false,