【问题标题】:Python Sockets multiple file transfersPython Sockets 多文件传输
【发布时间】:2021-11-15 17:13:47
【问题描述】:

所以我有一个 python 程序,它基本上让客户端连接到服务器并向其发送一个 excel 文件,该文件用作优化问题的输入。然后我希望服务器将此优化的输出(也是一个 excel 文件)发送回客户端。模型本身需要大约一分钟才能解决,我认为这会导致客户端尝试“过早”接收输出时出现一些问题。

客户端代码:

SEPARATOR = "<SEPARATOR>"
BUFFER_SIZE = 4096
HEADER = 64
HEADERSIZE = 10
port = 1234
FORMAT = 'utf-8'
DISCONNECT_MESSAGE = "!DISCONNECT"
host = "123.45.678.910"

s = socket.socket()
s.connect((host, port))


filename = "input/Model_Input.xlsx"
filesize = os.path.getsize(filename)
s.send(f"{filename}{SEPARATOR}{filesize}".encode())

with open(filename, "rb") as f:
    while True:
        bytes_read = f.read(BUFFER_SIZE)
        if not bytes_read:
            break
        s.sendall(bytes_read)

out_received = s.recv(BUFFER_SIZE).decode()
out_filename, out_filesize = out_received.split(SEPARATOR)
out_filename = os.path.basename(out_filename)
out_filesize = int(out_filesize)

with open(out_filename, "wb") as h:
    while True:
        out_bytes_read = s.recv(BUFFER_SIZE)
        if not out_bytes_read:
            break
        h.write(out_bytes_read)

以及服务器代码:

SERVER_PORT = 1234
SERVER_HOST = socket.gethostbyname(socket.gethostname())
BUFFER_SIZE = 4096
SEPARATOR = "<SEPARATOR>"

s = socket.socket()
s.bind((SERVER_HOST, SERVER_PORT))

s.listen(5)
client_socket, address = s.accept()

received = client_socket.recv(BUFFER_SIZE).decode()
filename, filesize = received.split(SEPARATOR)
filename = os.path.basename(filename)
filesize = int(filesize)

with open(filename, "wb") as f:
    while True:
        bytes_read = client_socket.recv(BUFFER_SIZE)
        if not bytes_read:
            break
        f.write(bytes_read)


##################
##  MODEL CODE  ##
##################

outfilename = 'Model_Output.xlsx'
outfilesize = os.path.getsize(filename)
client_socket.send(f"{outfilename}{SEPARATOR}{outfilesize}".encode())

with open(outfilename, "rb") as h:
    while True:
        # read the bytes from the file
        bytes_readed = h.read(BUFFER_SIZE)
        if not bytes_readed:
            break
        client_socket.sendall(bytes_readed)

我能够将输入文件发送到服务器并让模型运行,并将输出保存到存储中。但是,一旦我添加部分以尝试将其发送回客户端,它就会卡住。它仍然成功地发送和接收输入文件,但是模型永远不会运行。客户端和服务器都没有断开连接,它们似乎都被卡住了。

谢谢

【问题讨论】:

  • 套接字是零保证的数据流,当您发送(比如说)4096 个字节时,接收端将在一个、两个或十三个或 N 个块中接收这 4096 个字节;您必须包含一个协议以确保通信正常;例如,如果 碰巧被接收为 ),那么它将永远不会在另一端被识别。

标签: python file-transfer python-sockets


【解决方案1】:

对于某人(即我)来说,我可能很难远程调试这种类型的代码,因此我无法真正指出一定是问题所在的特定代码行。但是,如果您的客户端和服务器在同一台机器上运行,那么开始的客户端代码中可能存在问题:

with open(out_filename, "wb") as h:
    while True:
        out_bytes_read = s.recv(BUFFER_SIZE)
        if not out_bytes_read:
            break
        h.write(out_bytes_read)

open被执行时,这会将文件大小设置为0。服务器同时正在读取这个文件并将其传输到文件中,可能会发现现在只有0字节。但它已经发送了一个“标头”,表示有 N 个字节,其中 N 为非零。但这与您描述的问题不同。但它也可能发生在另一个方向。也就是说,当客户端正在发送文件并且服务器正在打开文件以进行输出时,它现在正在将客户端仍在读取的文件归零。下面的代码从两个方向解决了这个问题。当然,如果您的客户端和服务器位于不同的计算机上,而不会同时访问相同的文件,那么我所描述的就不是问题了。反正还没有。

但是,我可以提供一种稍微不同的方法,它似乎确实有效:

请参阅 Python 3 手册中 Socket Programming HOWTO 文章中的Using a Socket 部分。我已经采纳了使用固定长度消息的建议。这可能会更费力一点,但也更全面。这意味着如果要传输文件名,则必须首先将编码文件名的长度作为固定长度长度的消息传输(3 个字节可以处理长度最大为 999 字节的编码文件名),然后才能传输编码文件名。同样,我们将文件长度传输为 9 字节长度(左填充零),它可以处理最大为 999,999,999 字节的文件(我将宽度设置为 9 任意)。我有两个函数,receive_msgsend_msg,它们将稳健地发送和接收完整的字节消息,客户端和服务器都可以使用它们。这些是根据文章中的MySocket.mysendMySocket.myreceive 方法建模的。

我假设服务器在终止之前应该能够处理多个请求。实际上,它应该能够同时处理请求。为此,服务器将请求传递给线程池工作函数process_request 进行处理。目前尚不清楚所谓的“示范法典”的性质是什么。假设执行此计算的函数process_model 是 CPU 密集型的,process_request 将传递一个多处理池实例,该实例可用于执行 process_model 处理,这样 CPU 密集型处理部分就不会受到限制由全局解释器锁定。如果不涉及真正的 CPU 密集型处理,则删除创建多处理池的代码,然后将 process_model 作为常规函数调用。

服务器代码

import socket
from multiprocessing.pool import Pool, ThreadPool
import os.path

BUFFER_SIZE = 4096
SERVER_PORT = 1234
SERVER_HOST = socket.gethostbyname(socket.gethostname())

def receive_msg(sock, msg_length):
    chunks = []
    bytes_recd = 0
    while bytes_recd < msg_length:
        chunk = sock.recv(min(msg_length - bytes_recd, BUFFER_SIZE))
        if chunk == b'':
            raise RuntimeError("socket connection broken")
        chunks.append(chunk)
        bytes_recd = bytes_recd + len(chunk)
    return b''.join(chunks)

def send_msg(sock, msg):
    msg_length = len(msg)
    totalsent = 0
    while totalsent < msg_length:
        sent = sock.send(msg[totalsent:])
        if sent == 0:
            raise RuntimeError("socket connection broken")
        totalsent = totalsent + sent

def server():
    s = socket.socket()
    s.bind((SERVER_HOST, SERVER_PORT))
    s.listen(5)
    process_pool = Pool(5)
    thread_pool = ThreadPool(5)
    while True:
        client_socket, address = s.accept()
        thread_pool.apply_async(process_request, args=(process_pool, client_socket))

def process_request(process_pool, s):
    # Fixed length fields:
    # width 3 for filename length, followed by filename, width 9 for filesize
    filename_size = int(receive_msg(s, 3).decode())
    filename = receive_msg(s, filename_size).decode()
    filename = os.path.basename(filename)
    filesize = int(receive_msg(s, 9).decode())

    msg = receive_msg(s, filesize)
    with open(filename, "wb") as f:
        f.write(msg)

    # Assuming processing the model is CPU-intensive,
    # we use a process pool for doing that:
    out_filename = process_pool.apply(process_model, args=(filename,))
    out_filesize = os.path.getsize(out_filename)
    encoded_filename = out_filename.encode()
    msg1 = b"%03d%s%09dfilesize" % (len(encoded_filename), encoded_filename, out_filesize)
    with open(out_filename, "rb") as h:
        msg2 = h.read()
    send_msg(s, msg1)
    send_msg(s, msg2)

def process_model(filename):
    ...
    # Returned filename should probably be a function of the passed filename
    return 'Model_Output.xlsx' # name of the output file

if __name__ == '__main__':
    server()

客户端代码

import socket
import os.path

BUFFER_SIZE = 4096
port = 1234
host = "123.45.678.910"

def receive_msg(sock, msg_length):
    chunks = []
    bytes_recd = 0
    while bytes_recd < msg_length:
        chunk = sock.recv(min(msg_length - bytes_recd, BUFFER_SIZE))
        if chunk == b'':
            raise RuntimeError("socket connection broken")
        chunks.append(chunk)
        bytes_recd = bytes_recd + len(chunk)
    return b''.join(chunks)

def send_msg(sock, msg):
    msg_length = len(msg)
    totalsent = 0
    while totalsent < msg_length:
        sent = sock.send(msg[totalsent:])
        if sent == 0:
            raise RuntimeError("socket connection broken")
        totalsent = totalsent + sent

def client():
    s = socket.socket()
    s.connect((host, port))

    filename = "input/Model_Input.xlsx"
    filesize = os.path.getsize(filename)
    # Fixed length fields:
    # width 3 for filename length, followed by filename, width 9 for filesize
    encoded_filename = filename.encode()
    msg1 = b"%03d%s%09dfilesize" % (len(encoded_filename), encoded_filename, filesize)
    with open(filename, "rb") as f:
        msg2 = f.read()
    send_msg(s, msg1)
    send_msg(s, msg2)

    out_filename_size = int(receive_msg(s, 3).decode())
    out_filename = receive_msg(s, out_filename_size).decode()
    out_filename = os.path.basename(out_filename)
    out_filesize = int(receive_msg(s, 9).decode())

    msg = receive_msg(s, out_filesize)
    with open(out_filename, "wb") as h:
        h.write(msg)

if __name__ == '__main__':
    client()

更新

通过将服务实现为远程过程调用,可以大大简化整个编程。代码基于Python Cookbook, 3rd Edition

服务器

import socket
import pickle
from multiprocessing.connection import Listener
from threading import Thread
from multiprocessing.pool import Pool
import os.path

SERVER_PORT = 1234
SERVER_HOST = socket.gethostbyname(socket.gethostname())

class RPCHandler:
    def __init__(self):
        self._functions = { }

    def register_function(self, func):
        self._functions[func.__name__] = func

    def handle_connection(self, connection):
        try:
            while True:
                # Receive a message
                func_name, args, kwargs = pickle.loads(connection.recv())
                # Run the RPC and send a response
                try:
                    r = self._functions[func_name](*args,**kwargs)
                    connection.send(pickle.dumps(r))
                except Exception as e:
                    connection.send(pickle.dumps(e))
        except EOFError:
            pass

def server():
    global process_pool

    handler = RPCHandler()
    handler.register_function(process_request)
    sock = Listener((SERVER_HOST, SERVER_PORT))
    process_pool = Pool(5)
    while True:
        client = sock.accept()
        t = Thread(target=handler.handle_connection, args=(client,))
        t.daemon = True
        t.start()

def process_request(filename, contents):
    filename = os.path.basename(filename)
    with open(filename, "wb") as f:
        f.write(contents)

    # Assuming processing the model is CPU-intensive,
    # we use a process pool for doing that:
    out_filename = process_pool.apply(process_model, args=(filename,))
    with open(out_filename, "rb") as h:
        out_contents = h.read()
    return (out_filename, out_contents)

def process_model(filename):
    ...
    # Returned filename should probably be a function of the passed filename
    return 'Model_Output.xlsx' # name of the output file

if __name__ == '__main__':
    server()

客户

import pickle
import socket
import os.path
from multiprocessing.connection import Client

port = 1234
host = "123.45.678.910"

class RPCProxy:
    def __init__(self, connection):
        self._connection = connection

    def __getattr__(self, name):
        def do_rpc(*args, **kwargs):
            self._connection.send(pickle.dumps((name, args, kwargs)))
            result = pickle.loads(self._connection.recv())
            if isinstance(result, Exception):
                raise result
            return result
        return do_rpc

def client():
    c = Client((host, port))
    proxy = RPCProxy(c)

    filename = "input/Model_Input.xlsx"
    with open(filename, "rb") as f:
        contents = f.read()
    (out_filename, out_contents) = proxy.process_request(filename, contents)
    out_filename = os.path.basename(out_filename)
    with open(out_filename, "wb") as h:
        h.write(out_contents)

if __name__ == '__main__':
    client()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-12-12
    • 2021-08-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多