【问题标题】:Handle multiple messages with Queue get()使用 Queue get() 处理多条消息
【发布时间】:2014-12-15 23:14:34
【问题描述】:

感谢@user5402 之前的solution

我正在尝试处理多条排队的消息。代码如下:

import sys
import socket
from multiprocessing import Process, Queue

UDP_ADDR = ("", 13000)

def send(m):
    sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) 
    sock.sendto(m, UDP_ADDR)

def receive(q):
    buf = 1024
    Sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
    Sock.bind(UDP_ADDR)
    while True:
        (data, addr) = Sock.recvfrom(buf)
        q.put(data)

在客户端函数中,我想处理多个具有连锁影响的消息。

def client():
    q = Queue()
    r = Process(target = receive, args=(q,))
    r.start()

    print "client loop started"
    while True:
        m = q.get()
        print "got:", m
        while m == "start":
            print "started"
            z = q.get()
            if z == "stop":
                return
    print "loop ended"
    r.terminate()

所以当start 被发送时,它会进入一个无限打印"started" 的while 循环,并等待stop 消息通过。上面的client 代码不起作用。

下面是启动函数的代码:

if __name__ == '__main__':
    args = sys.argv
    if len(args) > 1:
        send(args[1])
    else:
        client()

【问题讨论】:

    标签: python multithreading sockets queue python-multithreading


    【解决方案1】:

    你可以这样写客户端循环:

    print "client loop started"
    while True:
        m = q.get()
        print "waiting for start, got:", m
        if m == "start":
          while True:
            try:
              m = q.get(False)
            except:
              m = None
            print "waiting for stop, got:", m
            if m == "stop":
              break
    

    根据您的 cmets,这将是一个更好的方法:

    import sys
    import socket
    import Queue as Q
    import time
    from multiprocessing import Process, Queue
    
    UDP_ADDR = ("", 13000)
    
    def send(m):
        sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) 
        sock.sendto(m, UDP_ADDR)
    
    def receive(q):
        buf = 1024
        Sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
        Sock.bind(UDP_ADDR)
        while True:
          (data, addr) = Sock.recvfrom(buf)
          q.put(data)
    
    def doit():
      # ... what the processing thread will do ...
      while True:
        print "sleeping..."
        time.sleep(3)
    
    def client():
      q = Queue()
      r = Process(target = receive, args=(q,))
      r.start()
    
      print "client loop started"
      t = None   # the processing thread
      while True:
          m = q.get()
          if m == "start":
            if t:
              print "processing thread already started"
            else:
              t = Process(target = doit)
              t.start()
              print "processing thread started"
          elif m == "stop":
            if t:
              t.terminate()
              t = None
              print "processing thread stopped"
            else:
              print "processing thread not running"
          elif m == "quit":
            print "shutting down"
            if t:
              t.terminate()
              t = None  # play it safe
            break
          else:
            print "huh?"
      r.terminate()
    
    if __name__ == '__main__':
      args = sys.argv
      if len(args) > 1:
        send(args[1])
      else:
        client()
    

    【讨论】:

    • 这不起作用,print "waiting for stop, got:", m 只返回一个,什么时候应该一直打印,直到发送stop。它一直打印的原因是因为我将插入更多需要循环的代码。 @user5402
    • 答案已更新 - 但您不想这样做。
    • 我正在创建的程序,它不断地使用蓝牙和 bluez 搜索信号强度。我知道这是不好的做法,但这是唯一可行的方法。非常感谢你的帮助! @user5402
    • 这实际上不起作用,它返回一个error: [Errno 98] Address already in use@user5402
    • 向我展示你的完整程序。您有另一个程序/进程正在侦听同一端口。
    猜你喜欢
    • 1970-01-01
    • 2014-12-27
    • 2018-01-04
    • 1970-01-01
    • 2018-07-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-10-03
    相关资源
    最近更新 更多