aboutsummaryrefslogtreecommitdiff
path: root/weed/messaging/client/subscriber.go
diff options
context:
space:
mode:
authorChris Lu <chrislusf@users.noreply.github.com>2020-05-17 17:39:16 -0700
committerGitHub <noreply@github.com>2020-05-17 17:39:16 -0700
commite0e31e67a809d00c99edaa299531c7ce4d4750dc (patch)
tree0f890277ef14c748faed4fecb7f8b8d4edeb9849 /weed/messaging/client/subscriber.go
parentb4e02ec525a6ec87b26686202307896faf3296a7 (diff)
parent081ee6fe349b519da8ea54cf3cdc17d2b15c5a71 (diff)
downloadseaweedfs-e0e31e67a809d00c99edaa299531c7ce4d4750dc.tar.xz
seaweedfs-e0e31e67a809d00c99edaa299531c7ce4d4750dc.zip
Merge pull request #1318 from chrislusf/msg_channel
Add messaging, add channel
Diffstat (limited to 'weed/messaging/client/subscriber.go')
-rw-r--r--weed/messaging/client/subscriber.go91
1 files changed, 0 insertions, 91 deletions
diff --git a/weed/messaging/client/subscriber.go b/weed/messaging/client/subscriber.go
deleted file mode 100644
index 53e7ffc7d..000000000
--- a/weed/messaging/client/subscriber.go
+++ /dev/null
@@ -1,91 +0,0 @@
-package client
-
-import (
- "context"
- "io"
- "time"
-
- "github.com/chrislusf/seaweedfs/weed/pb/messaging_pb"
-)
-
-type Subscriber struct {
- subscriberClients []messaging_pb.SeaweedMessaging_SubscribeClient
- subscriberId string
-}
-
-func (mc *MessagingClient) NewSubscriber(subscriberId, namespace, topic string, startTime time.Time) (*Subscriber, error) {
- // read topic configuration
- topicConfiguration := &messaging_pb.TopicConfiguration{
- PartitionCount: 4,
- }
- subscriberClients := make([]messaging_pb.SeaweedMessaging_SubscribeClient, topicConfiguration.PartitionCount)
-
- for i := 0; i < int(topicConfiguration.PartitionCount); i++ {
- client, err := mc.setupSubscriberClient(subscriberId, namespace, topic, int32(i), startTime)
- if err != nil {
- return nil, err
- }
- subscriberClients[i] = client
- }
-
- return &Subscriber{
- subscriberClients: subscriberClients,
- subscriberId: subscriberId,
- }, nil
-}
-
-func (mc *MessagingClient) setupSubscriberClient(subscriberId, namespace, topic string, partition int32, startTime time.Time) (messaging_pb.SeaweedMessaging_SubscribeClient, error) {
-
- stream, err := messaging_pb.NewSeaweedMessagingClient(mc.grpcConnection).Subscribe(context.Background())
- if err != nil {
- return nil, err
- }
-
- // send init message
- err = stream.Send(&messaging_pb.SubscriberMessage{
- Init: &messaging_pb.SubscriberMessage_InitMessage{
- Namespace: namespace,
- Topic: topic,
- Partition: partition,
- StartPosition: messaging_pb.SubscriberMessage_InitMessage_TIMESTAMP,
- TimestampNs: startTime.UnixNano(),
- SubscriberId: subscriberId,
- },
- })
- if err != nil {
- return nil, err
- }
-
- // process init response
- initResponse, err := stream.Recv()
- if err != nil {
- return nil, err
- }
- if initResponse.Redirect != nil {
- // TODO follow redirection
- }
-
- return stream, nil
-
-}
-
-func (s *Subscriber) doSubscribe(partition int, processFn func(m *messaging_pb.Message)) error {
- for {
- resp, listenErr := s.subscriberClients[partition].Recv()
- if listenErr == io.EOF {
- return nil
- }
- if listenErr != nil {
- println(listenErr.Error())
- return listenErr
- }
- processFn(resp.Data)
- }
-}
-
-// 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++ {
- go s.doSubscribe(i, processFn)
- }
-}