使用单一的等待参数
不读取数据
This commit is contained in:
@@ -30,7 +30,7 @@ const (
|
|||||||
pingPeriod = (pongWait * 9) / 10
|
pingPeriod = (pongWait * 9) / 10
|
||||||
|
|
||||||
// Maximum message size allowed from peer.
|
// Maximum message size allowed from peer.
|
||||||
maxMessageSize = 512
|
// maxMessageSize = 512
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -90,7 +90,7 @@ func (c *Client) readPump() {
|
|||||||
c.conn.Close()
|
c.conn.Close()
|
||||||
log.WithField("client", c).Warn("exit")
|
log.WithField("client", c).Warn("exit")
|
||||||
}()
|
}()
|
||||||
c.conn.SetReadLimit(maxMessageSize)
|
// c.conn.SetReadLimit(maxMessageSize)
|
||||||
c.conn.SetReadDeadline(time.Now().Add(pongWait))
|
c.conn.SetReadDeadline(time.Now().Add(pongWait))
|
||||||
c.conn.SetPongHandler(func(string) error { c.conn.SetReadDeadline(time.Now().Add(pongWait)); return nil })
|
c.conn.SetPongHandler(func(string) error { c.conn.SetReadDeadline(time.Now().Add(pongWait)); return nil })
|
||||||
|
|
||||||
@@ -231,8 +231,9 @@ func (h *Hub) read() {
|
|||||||
for h.running == true {
|
for h.running == true {
|
||||||
|
|
||||||
if h.reader == nil {
|
if h.reader == nil {
|
||||||
|
log.Info("waiting ... ")
|
||||||
|
time.Sleep(writeWait)
|
||||||
h.reader = create(h.topic, h.offset)
|
h.reader = create(h.topic, h.offset)
|
||||||
time.Sleep(5 * time.Second)
|
|
||||||
}
|
}
|
||||||
m, err := h.reader.ReadMessage(context.Background())
|
m, err := h.reader.ReadMessage(context.Background())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user