|
@ -2,7 +2,7 @@ package pub_client |
|
|
|
|
|
|
|
|
import ( |
|
|
import ( |
|
|
"github.com/rdleal/intervalst/interval" |
|
|
"github.com/rdleal/intervalst/interval" |
|
|
"github.com/seaweedfs/seaweedfs/weed/mq/broker" |
|
|
|
|
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/mq/balancer" |
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/mq_pb" |
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/mq_pb" |
|
|
"google.golang.org/grpc" |
|
|
"google.golang.org/grpc" |
|
|
"google.golang.org/grpc/credentials/insecure" |
|
|
"google.golang.org/grpc/credentials/insecure" |
|
@ -46,7 +46,7 @@ func (p *TopicPublisher) Connect(bootstrapBroker string) error { |
|
|
|
|
|
|
|
|
func (p *TopicPublisher) Shutdown() error { |
|
|
func (p *TopicPublisher) Shutdown() error { |
|
|
|
|
|
|
|
|
if clients, found := p.partition2Broker.AllIntersections(0, broker.MaxPartitionCount); found { |
|
|
|
|
|
|
|
|
if clients, found := p.partition2Broker.AllIntersections(0, balancer.MaxPartitionCount); found { |
|
|
for _, client := range clients { |
|
|
for _, client := range clients { |
|
|
client.CloseSend() |
|
|
client.CloseSend() |
|
|
} |
|
|
} |
|
|