Comments (2)
You can create the Producer (using NewProducer()
) during the initialization of your service, store it to a global variable, and re-use it from all request goroutines. It is "thread-safe".
from go-nsq.
Is it possible to use singleton NSQ Producer? I'm running an API with NSQ Producer and it seems like i have to open and close connection everytime i hit endpoint.
Here is a example you can use:
package nq
import (
"encoding/json"
"errors"
"fmt"
"github.com/nsqio/go-nsq"
"math/rand"
"sync"
"time"
)
type Client struct {
Producer []*nsq.Producer
PLen int
topics map[string]struct{}
mux sync.Mutex
Prefix string
NsqdAddress []string
LookupdAddress []string
}
func New(prefix string, nsqdAddress []string, lookupdAddress []string) (*Client, error) {
if len(nsqdAddress) == 0 || len(lookupdAddress) == 0 {
return nil, errors.New("config invalid")
}
rand.Seed(time.Now().Unix())
c := &Client{
Producer: nil,
NsqdAddress: nsqdAddress,
LookupdAddress: lookupdAddress,
Prefix: prefix,
PLen: len(nsqdAddress),
topics: map[string]struct{}{},
}
for _, v := range nsqdAddress {
producer, err := nsq.NewProducer(v, nsq.NewConfig())
if err != nil {
return nil, err
}
c.Producer = append(c.Producer, producer)
}
return c, nil
}
func (c *Client) Pub(topic string, object interface{}) error {
if c.Prefix != "" {
topic = fmt.Sprintf("%s.%s", c.Prefix, topic)
}
body, err := json.Marshal(object)
if err != nil {
return err
}
return c.Producer[rand.Intn(c.PLen)].Publish(topic, body)
}
func (c *Client) DeferredPub(topic string, delay time.Duration, object interface{}) error {
if c.Prefix != "" {
topic = fmt.Sprintf("%s.%s", c.Prefix, topic)
}
body, err := json.Marshal(object)
if err != nil {
return err
}
return c.Producer[rand.Intn(c.PLen)].DeferredPublish(topic, delay, body)
}
func (c *Client) PubRaw(topic string, raw []byte) error {
if c.Prefix != "" {
topic = fmt.Sprintf("%s.%s", c.Prefix, topic)
}
return c.Producer[rand.Intn(c.PLen)].Publish(topic, raw)
}
func (c *Client) DeferredPubRaw(topic string, delay time.Duration, raw []byte) error {
if c.Prefix != "" {
topic = fmt.Sprintf("%s.%s", c.Prefix, topic)
}
return c.Producer[rand.Intn(c.PLen)].DeferredPublish(topic, delay, raw)
}
func (c *Client) Sub(topic, channel string, handler nsq.Handler) (err error) {
c.mux.Lock()
defer c.mux.Unlock()
if _, ok := c.topics[topic]; ok {
return nil
} else {
c.topics[topic] = struct{}{}
}
if c.Prefix != "" {
topic = fmt.Sprintf("%s.%s", c.Prefix, topic)
}
consumer, err := nsq.NewConsumer(topic, channel, nsq.NewConfig())
if err != nil {
return err
}
consumer.AddHandler(handler)
if err := consumer.ConnectToNSQLookupds(c.LookupdAddress); err != nil {
return err
}
return nil
}
func (c *Client) SubMany(topic, channel string, handler nsq.Handler, concurrency int) (err error) {
c.mux.Lock()
defer c.mux.Unlock()
if _, ok := c.topics[topic]; ok {
return nil
} else {
c.topics[topic] = struct{}{}
}
if c.Prefix != "" {
topic = fmt.Sprintf("%s.%s", c.Prefix, topic)
}
consumer, err := nsq.NewConsumer(topic, channel, nsq.NewConfig())
if err != nil {
return err
}
consumer.AddConcurrentHandlers(handler, concurrency)
if err := consumer.ConnectToNSQLookupds(c.LookupdAddress); err != nil {
return err
}
return nil
}
Test it:
func TestNew(t *testing.T) {
a, err := New("En-Shop", []string{"10.0.5.2:5150","10.0.5.3:5150","10.0.5.4:5150"}, []string{"10.0.5.2:5161","10.0.5.3:5161","10.0.5.4:5161"})
if err != nil {
fmt.Println(err.Error())
return
}
err = a.DeferredPubRaw("Test", time.Minute*10, []byte("ssss"))
if err != nil {
fmt.Println(err.Error())
return
}
}
from go-nsq.
Related Issues (20)
- panic: send closed channel - incomingMessages HOT 1
- feat: NSQ cluster or public network message synchronization
- v1.2.1 use DeferredPublish publish delay message,delay invalidation HOT 4
- Can you tag v1.0.9 HOT 1
- How to use Command? HOT 1
- Support NSQ in ArgoLabs Dataflow HOT 1
- the tag of v1.0.8, and the file of doc.go line 65, the value name p is error, in fact, it is producer. HOT 2
- Disconnects seem to be quite ungraceful HOT 2
- "go test" fails with "connection refused" HOT 2
- Client doesn't know it's disconnected HOT 2
- Panic in producer's router HOT 2
- Conn cannot execute multiple subscribe commands HOT 1
- error connecting to nsqd - dial tcp 127.0.0.1:4150: connect: connection refused? HOT 1
- An Example of multit-hreaded nsq consuming messages HOT 1
- producer: HTTP response to IDENTIFY causes memory allocation of over 1GB HOT 2
- how to clear data in nsq HOT 1
- MaxAttempts Callback HOT 1
- *: context.Context and timeouts
- How can I limit the queue size?
Recommend Projects
-
React
A declarative, efficient, and flexible JavaScript library for building user interfaces.
-
Vue.js
🖖 Vue.js is a progressive, incrementally-adoptable JavaScript framework for building UI on the web.
-
Typescript
TypeScript is a superset of JavaScript that compiles to clean JavaScript output.
-
TensorFlow
An Open Source Machine Learning Framework for Everyone
-
Django
The Web framework for perfectionists with deadlines.
-
Laravel
A PHP framework for web artisans
-
D3
Bring data to life with SVG, Canvas and HTML. 📊📈🎉
-
Recommend Topics
-
javascript
JavaScript (JS) is a lightweight interpreted programming language with first-class functions.
-
web
Some thing interesting about web. New door for the world.
-
server
A server is a program made to process requests and deliver data to clients.
-
Machine learning
Machine learning is a way of modeling and interpreting data that allows a piece of software to respond intelligently.
-
Visualization
Some thing interesting about visualization, use data art
-
Game
Some thing interesting about game, make everyone happy.
Recommend Org
-
Facebook
We are working to build community through open source technology. NB: members must have two-factor auth.
-
Microsoft
Open source projects and samples from Microsoft.
-
Google
Google ❤️ Open Source for everyone.
-
Alibaba
Alibaba Open Source for everyone
-
D3
Data-Driven Documents codes.
-
Tencent
China tencent open source team.
from go-nsq.