aboutsummaryrefslogtreecommitdiff
path: root/weed/msgqueue
diff options
context:
space:
mode:
authorChris Lu <chris.lu@gmail.com>2018-09-16 01:18:30 -0700
committerChris Lu <chris.lu@gmail.com>2018-09-16 01:18:30 -0700
commitd923ba2206a8238977e740f98f061242d61ed5e6 (patch)
treea2fcaf40d5f2398b2b8ad3c2a77a22579597cc1c /weed/msgqueue
parentbea4f6ca14615de23ba81a84cf61ba0f2a28b7f6 (diff)
downloadseaweedfs-d923ba2206a8238977e740f98f061242d61ed5e6.tar.xz
seaweedfs-d923ba2206a8238977e740f98f061242d61ed5e6.zip
renaming msgqueue to notification
Diffstat (limited to 'weed/msgqueue')
-rw-r--r--weed/msgqueue/configuration.go33
-rw-r--r--weed/msgqueue/kafka/kafka_queue.go79
-rw-r--r--weed/msgqueue/log/log_queue.go29
-rw-r--r--weed/msgqueue/message_queue.go14
4 files changed, 0 insertions, 155 deletions
diff --git a/weed/msgqueue/configuration.go b/weed/msgqueue/configuration.go
deleted file mode 100644
index 769c11835..000000000
--- a/weed/msgqueue/configuration.go
+++ /dev/null
@@ -1,33 +0,0 @@
-package msgqueue
-
-import (
- "github.com/chrislusf/seaweedfs/weed/glog"
- "github.com/spf13/viper"
-)
-
-var (
- MessageQueues []MessageQueue
-
- Queue MessageQueue
-)
-
-func LoadConfiguration(config *viper.Viper) {
-
- if config == nil {
- return
- }
-
- for _, store := range MessageQueues {
- if config.GetBool(store.GetName() + ".enabled") {
- viperSub := config.Sub(store.GetName())
- if err := store.Initialize(viperSub); err != nil {
- glog.Fatalf("Failed to initialize store for %s: %+v",
- store.GetName(), err)
- }
- Queue = store
- glog.V(0).Infof("Configure message queue for %s", store.GetName())
- return
- }
- }
-
-}
diff --git a/weed/msgqueue/kafka/kafka_queue.go b/weed/msgqueue/kafka/kafka_queue.go
deleted file mode 100644
index f373395f4..000000000
--- a/weed/msgqueue/kafka/kafka_queue.go
+++ /dev/null
@@ -1,79 +0,0 @@
-package kafka
-
-import (
- "github.com/Shopify/sarama"
- "github.com/chrislusf/seaweedfs/weed/glog"
- "github.com/chrislusf/seaweedfs/weed/msgqueue"
- "github.com/chrislusf/seaweedfs/weed/util"
- "github.com/golang/protobuf/proto"
-)
-
-func init() {
- msgqueue.MessageQueues = append(msgqueue.MessageQueues, &KafkaQueue{})
-}
-
-type KafkaQueue struct {
- topic string
- producer sarama.AsyncProducer
-}
-
-func (k *KafkaQueue) GetName() string {
- return "kafka"
-}
-
-func (k *KafkaQueue) Initialize(configuration util.Configuration) (err error) {
- glog.V(0).Infof("filer.msgqueue.kafka.hosts: %v\n", configuration.GetStringSlice("hosts"))
- glog.V(0).Infof("filer.msgqueue.kafka.topic: %v\n", configuration.GetString("topic"))
- return k.initialize(
- configuration.GetStringSlice("hosts"),
- configuration.GetString("topic"),
- )
-}
-
-func (k *KafkaQueue) initialize(hosts []string, topic string) (err error) {
- config := sarama.NewConfig()
- config.Producer.RequiredAcks = sarama.WaitForLocal
- config.Producer.Partitioner = sarama.NewHashPartitioner
- config.Producer.Return.Successes = true
- config.Producer.Return.Errors = true
- k.producer, err = sarama.NewAsyncProducer(hosts, config)
- k.topic = topic
- go k.handleSuccess()
- go k.handleError()
- return nil
-}
-
-func (k *KafkaQueue) SendMessage(key string, message proto.Message) (err error) {
- bytes, err := proto.Marshal(message)
- if err != nil {
- return
- }
-
- msg := &sarama.ProducerMessage{
- Topic: k.topic,
- Key: sarama.StringEncoder(key),
- Value: sarama.ByteEncoder(bytes),
- }
-
- k.producer.Input() <- msg
-
- return nil
-}
-
-func (k *KafkaQueue) handleSuccess() {
- for {
- pm := <-k.producer.Successes()
- if pm != nil {
- glog.V(3).Infof("producer message success, partition:%d offset:%d key:%v", pm.Partition, pm.Offset, pm.Key)
- }
- }
-}
-
-func (k *KafkaQueue) handleError() {
- for {
- err := <-k.producer.Errors()
- if err != nil {
- glog.Errorf("producer message error, partition:%d offset:%d key:%v valus:%s error(%v) topic:%s", err.Msg.Partition, err.Msg.Offset, err.Msg.Key, err.Msg.Value, err.Err, k.topic)
- }
- }
-}
diff --git a/weed/msgqueue/log/log_queue.go b/weed/msgqueue/log/log_queue.go
deleted file mode 100644
index d291d629c..000000000
--- a/weed/msgqueue/log/log_queue.go
+++ /dev/null
@@ -1,29 +0,0 @@
-package kafka
-
-import (
- "github.com/chrislusf/seaweedfs/weed/glog"
- "github.com/chrislusf/seaweedfs/weed/msgqueue"
- "github.com/chrislusf/seaweedfs/weed/util"
- "github.com/golang/protobuf/proto"
-)
-
-func init() {
- msgqueue.MessageQueues = append(msgqueue.MessageQueues, &LogQueue{})
-}
-
-type LogQueue struct {
-}
-
-func (k *LogQueue) GetName() string {
- return "log"
-}
-
-func (k *LogQueue) Initialize(configuration util.Configuration) (err error) {
- return nil
-}
-
-func (k *LogQueue) SendMessage(key string, message proto.Message) (err error) {
-
- glog.V(0).Infof("%v: %+v", key, message)
- return nil
-}
diff --git a/weed/msgqueue/message_queue.go b/weed/msgqueue/message_queue.go
deleted file mode 100644
index 3e2688698..000000000
--- a/weed/msgqueue/message_queue.go
+++ /dev/null
@@ -1,14 +0,0 @@
-package msgqueue
-
-import (
- "github.com/chrislusf/seaweedfs/weed/util"
- "github.com/golang/protobuf/proto"
-)
-
-type MessageQueue interface {
- // GetName gets the name to locate the configuration in message_queue.toml file
- GetName() string
- // Initialize initializes the file store
- Initialize(configuration util.Configuration) error
- SendMessage(key string, message proto.Message) error
-}