diff options
Diffstat (limited to 'weed/messaging')
| -rw-r--r-- | weed/messaging/broker/broker_append.go | 6 | ||||
| -rw-r--r-- | weed/messaging/client/subscriber.go | 2 |
2 files changed, 7 insertions, 1 deletions
diff --git a/weed/messaging/broker/broker_append.go b/weed/messaging/broker/broker_append.go index 7194dfcfc..26f24f4d3 100644 --- a/weed/messaging/broker/broker_append.go +++ b/weed/messaging/broker/broker_append.go @@ -98,6 +98,8 @@ func (broker *MessageBroker) assignAndUpload(topicConfig *messaging_pb.TopicConf return assignResult, uploadResult, nil } +var _ = filer_pb.FilerClient(&MessageBroker{}) + func (broker *MessageBroker) WithFilerClient(fn func(filer_pb.SeaweedFilerClient) error) (err error) { for _, filer := range broker.option.Filers { @@ -111,3 +113,7 @@ func (broker *MessageBroker) WithFilerClient(fn func(filer_pb.SeaweedFilerClient return } + +func (broker *MessageBroker) AdjustedUrl(hostAndPort string) string { + return hostAndPort +} diff --git a/weed/messaging/client/subscriber.go b/weed/messaging/client/subscriber.go index ddf1f82e6..2ebad4ce6 100644 --- a/weed/messaging/client/subscriber.go +++ b/weed/messaging/client/subscriber.go @@ -85,7 +85,7 @@ func (s *Subscriber) doSubscribe(partition int, processFn func(m *messaging_pb.M // Subscribe starts goroutines to process the messages func (s *Subscriber) Subscribe(processFn func(m *messaging_pb.Message)) { - for i:=0;i<len(s.subscriberClients);i++{ + for i := 0; i < len(s.subscriberClients); i++ { go s.doSubscribe(i, processFn) } } |
