mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2024-11-28 21:39:02 +08:00
Merge pull request #1794 from kmlebedev/recoveringRabbitMQ
This commit is contained in:
commit
755b524814
@ -113,6 +113,10 @@ func (k *GoCDKPubSubInput) Initialize(configuration util.Configuration, prefix s
|
||||
func (k *GoCDKPubSubInput) ReceiveMessage() (key string, message *filer_pb.EventNotification, onSuccessFn func(), onFailureFn func(), err error) {
|
||||
msg, err := k.sub.Receive(context.Background())
|
||||
if err != nil {
|
||||
var conn *amqp.Connection
|
||||
if k.sub.As(&conn) && conn.IsClosed() {
|
||||
glog.Fatalln(err)
|
||||
}
|
||||
return
|
||||
}
|
||||
onFailureFn = func() {
|
||||
|
Loading…
Reference in New Issue
Block a user