【问题标题】:Thread synchronization in PythonPython中的线程同步
【发布时间】:2012-03-20 05:41:24
【问题描述】:

我目前正在从事一个学校项目,其中的任务是设置一个线程服务器/客户端系统。系统中的每个客户端在连接到服务器时都应该在服务器上分配自己的线程。此外,我希望服务器运行其他线程,一个涉及来自命令行的输入,另一个涉及向所有客户端广播消息。但是,我无法让它按我的意愿运行。似乎线程正在相互阻塞。我希望我的程序在服务器侦听连接的客户端的“同时”从命令行获取输入,依此类推。

我是 python 编程和多线程的新手,虽然我认为我的想法很好,但我并不惊讶我的代码不起作用。问题是我不确定我将如何实现不同线程之间的消息传递。我也不确定如何正确实施资源锁定命令。我将在这里发布我的服务器文件和客户端文件的代码,希望有人可以帮助我。我认为这实际上应该是两个相对简单的脚本。我试图在某种程度上尽可能好地评论我的代码。

import select
import socket
import sys
import threading
import client

class Server:

#initializing server socket
def __init__(self, event):
    self.host = 'localhost'
    self.port = 50000
    self.backlog = 5
    self.size = 1024
    self.server = None
    self.server_running = False
    self.listen_threads = []
    self.local_threads = []
    self.clients = []
    self.serverSocketLock = None
    self.cmdLock = None
    #here i have also declared some events for the command line input
    #and the receive function respectively, not sure if correct
    self.cmd_event = event
    self.socket_event = event

def openSocket(self):
    #binding server to port
    try: 
        self.server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        self.server.bind((self.host, self.port))
        self.server.listen(5)
        print "Listening to port " + str(self.port) + "..."
    except socket.error, (value,message):
        if self.server:
            self.server.close()
        print "Could not open socket: " + message
        sys.exit(1)

def run(self):
    self.openSocket()

    #making Rlocks for the socket and for the command line input

    self.serverSocketLock = threading.RLock()
    self.cmdLock = threading.RLock()

    #set blocking to non-blocking
    self.server.setblocking(0)
    #making two threads always running on the server,
    #one for the command line input, and one for broadcasting (sending)
    cmd_thread = threading.Thread(target=self.server_cmd)
    broadcast_thread = threading.Thread(target=self.broadcast,args=[self.clients])
    cmd_thread.daemon = True
    broadcast_thread.daemon = True
    #append the threads to thread list
    self.local_threads.append(cmd_thread)
    self.local_threads.append(broadcast_thread)

    cmd_thread.start()
    broadcast_thread.start()


    self.server_running = True
    while self.server_running:

        #connecting to "knocking" clients
        try:
            c = client.Client(self.server.accept())
            self.clients.append(c)
            print "Client " + str(c.address) + " connected"

            #making a thread for each clientn and appending it to client list
            listen_thread = threading.Thread(target=self.listenToClient,args=[c])
            self.listen_threads.append(listen_thread)
            listen_thread.daemon = True
            listen_thread.start()
            #setting event "client has connected"
            self.socket_event.set()

        except socket.error, (value, message):
            continue

    #close threads

    self.server.close()
    print "Closing client threads"
    for c in self.listen_threads:
        c.join()

def listenToClient(self, c):

    while self.server_running:

        #the idea here is to wait until the thread gets the message "client
        #has connected"
        self.socket_event.wait()
        #then clear the event immidiately...
        self.socket_event.clear()
        #and aquire the socket resource
        self.serverSocketLock.acquire()

        #the below is the receive thingy

        try:
            recvd_data = c.client.recv(self.size)
            if recvd_data == "" or recvd_data == "close\n":
                print "Client " + str(c.address) + (" disconnected...")
                self.socket_event.clear()
                self.serverSocketLock.release()
                return

            print recvd_data

            #I put these here to avoid locking the resource if no message 
            #has been received
            self.socket_event.clear()
            self.serverSocketLock.release()
        except socket.error, (value, message):
            continue            

def server_cmd(self):

    #this is a simple command line utility
    while self.server_running:

        #got to have a smart way to make this work

        self.cmd_event.wait()
        self.cmd_event.clear()
        self.cmdLock.acquire()


        cmd = sys.stdin.readline()
        if cmd == "":
            continue
        if cmd == "close\n":
            print "Server shutting down..."
            self.server_running = False

        self.cmdLock.release()


def broadcast(self, clients):
    while self.server_running:

        #this function will broadcast a message received from one
        #client, to all other clients, but i guess any thread
        #aspects applied to the above, will work here also

        try:
            send_data = sys.stdin.readline()
            if send_data == "":
                continue
            else:
                for c in clients:
                    c.client.send(send_data)
            self.serverSocketLock.release()
            self.cmdLock.release()
        except socket.error, (value, message):
            continue

if __name__ == "__main__":
e = threading.Event()
s = Server(e)
s.run()

然后是客户端文件

import select
import socket
import sys
import server
import threading

class Client(threading.Thread):

#initializing client socket

def __init__(self,(client,address)):
    threading.Thread.__init__(self) 
    self.client = client 
    self.address = address
    self.size = 1024
    self.client_running = False
    self.running_threads = []
    self.ClientSocketLock = None

def run(self):

    #connect to server
    self.client.connect(('localhost',50000))

    #making a lock for the socket resource
    self.clientSocketLock = threading.Lock()
    self.client.setblocking(0)
    self.client_running = True

    #making two threads, one for receiving messages from server...
    listen = threading.Thread(target=self.listenToServer)

    #...and one for sending messages to server
    speak = threading.Thread(target=self.speakToServer)

    #not actually sure wat daemon means
    listen.daemon = True
    speak.daemon = True

    #appending the threads to the thread-list
    self.running_threads.append(listen)
    self.running_threads.append(speak)
    listen.start()
    speak.start()

    #this while-loop is just for avoiding the script terminating
    while self.client_running:
        dummy = 1

    #closing the threads if the client goes down
    print "Client operating on its own"
    self.client.close()

    #close threads
    for t in self.running_threads:
        t.join()
    return

#defining "listen"-function
def listenToServer(self):
    while self.client_running:

        #here i acquire the socket to this function, but i realize I also
        #should have a message passing wait()-function or something
        #somewhere
        self.clientSocketLock.acquire()

        try:
            data_recvd = self.client.recv(self.size)
            print data_recvd
        except socket.error, (value,message):
            continue

        #releasing the socket resource
        self.clientSocketLock.release()

#defining "speak"-function, doing much the same as for the above function       
def speakToServer(self):
    while self.client_running:
        self.clientSocketLock.acquire()
        try:
            send_data = sys.stdin.readline()
            if send_data == "close\n":
                print "Disconnecting..."
                self.client_running = False
            else:
                self.client.send(send_data)
        except socket.error, (value,message):
            continue

        self.clientSocketLock.release()

if __name__ == "__main__":
c = Client((socket.socket(socket.AF_INET, socket.SOCK_STREAM),'localhost'))
c.run()

我知道这有很多代码行供您阅读,但正如我所说,我认为其中的概念和脚本本身应该很容易理解。如果有人可以帮助我以适当的方式同步我的线程,那将非常感激 =)

提前致谢

---编辑---

好的。所以我现在已经将我的代码简化为只在服务器和客户端模块中包含发送和接收功能。连接到服务器的客户端拥有自己的线程,两个模块中的发送和接收函数在各自独立的线程中运行。这就像一个魅力,服务器模块中的广播功能将它从一个客户端获取到所有客户端的字符串回显。到目前为止一切顺利!

我希望我的脚本做的下一件事是在客户端模块中执行特定命令,即“关闭”以关闭客户端,并加入线程列表中的所有正在运行的线程。我使用事件标志来通知 listenToServer 和主线程 speakToServer 线程已读取输入“关闭”。似乎主线程跳出了它的while循环并启动了应该加入其他线程的for循环。但在这里它挂起。即使在设置事件标志时 server_running 应该设置为 False,listenToServer 线程中的 while 循环似乎也永远不会停止。

我在这里只发布客户端模块,因为我想让这两个线程同步的答案也与在客户端和服务器模块中同步更多线程有关。

import select
import socket
import sys
import server_bygg0203
import threading
from time import sleep

class Client(threading.Thread):

#initializing client socket

def __init__(self,(client,address)):

threading.Thread.__init__(self) 
self.client = client 
self.address = address
self.size = 1024
self.client_running = False
self.running_threads = []
self.ClientSocketLock = None
self.disconnected = threading.Event()

def run(self):

#connect to server
self.client.connect(('localhost',50000))

#self.client.setblocking(0)
self.client_running = True

#making two threads, one for receiving messages from server...
listen = threading.Thread(target=self.listenToServer)

#...and one for sending messages to server
speak = threading.Thread(target=self.speakToServer)

#not actually sure what daemon means
listen.daemon = True
speak.daemon = True

#appending the threads to the thread-list
self.running_threads.append((listen,"listen"))
self.running_threads.append((speak, "speak"))
listen.start()
speak.start()

while self.client_running:

    #check if event is set, and if it is
    #set while statement to false

    if self.disconnected.isSet():
        self.client_running = False 

#closing the threads if the client goes down
print "Client operating on its own"
self.client.shutdown(1)
self.client.close()

#close threads

#the script hangs at the for-loop below, and
#refuses to close the listen-thread (and possibly
#also the speak thread, but it never gets that far)

for t in self.running_threads:
    print "Waiting for " + t[1] + " to close..."
    t[0].join()
self.disconnected.clear()
return

#defining "speak"-function      
def speakToServer(self):

#sends strings to server
while self.client_running:
    try:
        send_data = sys.stdin.readline()
        self.client.send(send_data)

        #I want the "close" command
        #to set an event flag, which is being read by all other threads,
        #and, at the same time set the while statement to false

        if send_data == "close\n":
            print "Disconnecting..."
            self.disconnected.set()
            self.client_running = False
    except socket.error, (value,message):
        continue
return

#defining "listen"-function
def listenToServer(self):

#receives strings from server
while self.client_running:

    #check if event is set, and if it is
    #set while statement to false

    if self.disconnected.isSet():
        self.client_running = False
    try:
        data_recvd = self.client.recv(self.size)
        print data_recvd
    except socket.error, (value,message):
        continue
return

if __name__ == "__main__":
c = Client((socket.socket(socket.AF_INET, socket.SOCK_STREAM),'localhost'))
c.run()

稍后,当我启动并运行这个服务器/客户端系统时,我将在我们实验室的一些电梯型号上使用这个系统,每个客户端都会接收楼层指令或“向上”和“向下”呼叫。服务器将运行分配算法并更新客户端上最适合请求订单的电梯队列。我意识到这是一个很长的路要走,但我想一个人应该只迈出一步 =)

希望有人有时间对此进行调查。提前致谢。

【问题讨论】:

  • 我认为您应该更明确地说明您的代码如何“不起作用”。你期待什么行为?你观察到什么行为? “不起作用”可能意味着很多事情。

标签: python multithreading sockets client


【解决方案1】:

备份你的工作,然后扔掉它 - 部分。

您需要分段实现您的程序,并在执行过程中测试每个部分。首先,处理程序的input 部分。不要担心如何广播您收到的输入。而是担心您是否能够通过套接字成功且重复地接收输入。到目前为止 - 非常好。

现在,我假设您希望通过向其他附加客户端广播来对此输入做出反应。太糟糕了,你还不能这样做!因为,我在上面的段落中留下了一个小细节。你必须设计一个协议。

什么是协议?这是一套沟通规则。您的服务器如何知道客户端何时完成发送数据?它是否被某些特殊字符终止?或者也许您将要发送的消息的大小编码为消息的第一个或两个字节。

结果证明这是很多工作,不是吗? :-)

什么是简单协议。 line-oriented 协议很简单。一次读取 1 个字符,直到到达记录终止符结尾 - '\n'。因此,客户端会将这样的记录发送到您的服务器--

直升机\n MSG DAVE 你的孩子在哪里?\n

所以,假设您设计了这个简单的协议,请实施它。现在,不要担心多线程的东西!只需担心它的工作。

您当前的协议是读取 1024 个字节。这可能还不错,只要确保从客户端发送 1024 字节的消息即可。

设置好协议内容后,继续对输入做出反应。但是现在你需要一些可以读取输入的东西。一旦完成,我们就可以担心用它做点什么了。

jdi 是对的,你有太多的程序要处理。碎片更容易修复。

【讨论】:

    【解决方案2】:

    我在这段代码中看到的最大问题是您有太多的事情要立即进行,无法轻松地调试您的问题。由于逻辑变得非常非线性,线程可能会变得非常复杂。尤其是当您不得不担心与锁同步时。

    您看到客户端相互阻塞的原因是您在服务器的 listenToClient() 循环中使用 serverSocketLock 的方式。老实说,这不完全是您的代码现在的问题,但是当我开始调试它并将套接字变成阻塞套接字时,它就成了问题。如果您将每个连接放入自己的线程并从中读取,那么这里没有理由使用全局服务器锁。它们都可以同时从自己的套接字中读取,这就是线程的目的。

    这是我给你的建议:

    1. 摆脱所有你不需要的锁和多余的线程,从头开始
    2. 让客户端像您一样连接,并像您一样将它们放入他们的线程中。并且只需让他们每秒发送一次数据。验证您可以让多个客户端连接和发送,并且您的服务器正在循环和接收。完成这一部分后,您可以继续进行下一部分。
    3. 现在您已将套接字设置为非阻塞。当数据没有准备好时,这导致它们都在循环中快速旋转。由于您是线程,因此您应该将它们设置为阻塞。然后阅读器线程将简单地坐下来等待数据并立即响应。

    当线程访问共享资源时使用锁。您显然需要在任何时候线程尝试修改服务器属性,如列表或值。但不是在他们使用自己的私有套接字时。

    您用来触发读者的事件在这里似乎没有必要。您已收到客户,然后您启动线程。所以它准备好了。

    简而言之...简化并一次测试一点。当它工作时,添加更多。现在线程和锁太多了。

    这是您的listenToClient 方法的简化示例:

    def listenToClient(self, c):
        while self.server_running:
            try:
                recvd_data = c.client.recv(self.size)
                print "received:", c, recvd_data
                if recvd_data == "" or recvd_data == "close\n":
                    print "Client " + str(c.address) + (" disconnected...")
                    return
    
                print recvd_data
    
            except socket.error, (value, message):
                if value == 35:
                    continue 
                else:
                    print "Error:", value, message  
    

    【讨论】:

    • 好的!似乎是一些合理的建议。当我有时间时,我会研究这个。非常感谢您花时间调试,当我启动并运行我的新代码时,我会发布答案 =)
    • 当然可以。如果您需要再次查看您的简化代码,请告诉我。
    猜你喜欢
    • 2012-05-01
    • 1970-01-01
    • 1970-01-01
    • 2011-11-09
    • 2014-12-15
    • 1970-01-01
    • 2011-12-28
    • 1970-01-01
    • 2014-04-30
    相关资源
    最近更新 更多