forked from mirror/redis
Create PubSub channel once
This commit is contained in:
parent
35248155a5
commit
15f14b8305
18
pubsub.go
18
pubsub.go
|
@ -29,6 +29,9 @@ type PubSub struct {
|
||||||
closed bool
|
closed bool
|
||||||
|
|
||||||
cmd *Cmd
|
cmd *Cmd
|
||||||
|
|
||||||
|
chOnce sync.Once
|
||||||
|
ch chan *Message
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *PubSub) conn() (*pool.Conn, error) {
|
func (c *PubSub) conn() (*pool.Conn, error) {
|
||||||
|
@ -346,10 +349,12 @@ func (c *PubSub) receiveMessage(timeout time.Duration) (*Message, error) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Channel returns a channel for concurrently receiving messages.
|
// Channel returns a Go channel for concurrently receiving messages.
|
||||||
// The channel is closed with PubSub.
|
// The channel is closed with PubSub. Receive or ReceiveMessage APIs
|
||||||
|
// can not be used after channel is created.
|
||||||
func (c *PubSub) Channel() <-chan *Message {
|
func (c *PubSub) Channel() <-chan *Message {
|
||||||
ch := make(chan *Message, 100)
|
c.chOnce.Do(func() {
|
||||||
|
c.ch = make(chan *Message, 100)
|
||||||
go func() {
|
go func() {
|
||||||
for {
|
for {
|
||||||
msg, err := c.ReceiveMessage()
|
msg, err := c.ReceiveMessage()
|
||||||
|
@ -359,11 +364,12 @@ func (c *PubSub) Channel() <-chan *Message {
|
||||||
}
|
}
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
ch <- msg
|
c.ch <- msg
|
||||||
}
|
}
|
||||||
close(ch)
|
close(c.ch)
|
||||||
}()
|
}()
|
||||||
return ch
|
})
|
||||||
|
return c.ch
|
||||||
}
|
}
|
||||||
|
|
||||||
func appendIfNotExists(ss []string, es ...string) []string {
|
func appendIfNotExists(ss []string, es ...string) []string {
|
||||||
|
|
Loading…
Reference in New Issue