You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Python Socket编程的TCP/IP通信架构实现问题咨询

Got it, let's walk through exactly how to build this communication flow with Python's socket and threading tools. I'll break it down into actionable parts with code examples you can adapt directly:

Core Architecture Overview

First, let's clarify each component's role to keep things aligned:

  • TCP Server: Acts as the central hub—manages connections from the controller and 3 worker clients, broadcasts controller messages to all workers, and forwards worker task results back to the controller.
  • Controller Client: Sends task instructions to the server and listens for aggregated results from workers.
  • Worker Clients: Receive task broadcasts from the server, execute their assigned tasks, and send completion data back to the server.
Step-by-Step Implementation

1. TCP Server Setup

The server needs to track connected clients (distinguishing controller vs workers) and route messages correctly. We'll use a thread per client to handle concurrent communication, plus a lock to ensure thread-safe access to shared connection data.

import socket
import threading

# Track connected clients: key = socket object, value = client type ('controller'/'worker')
client_sockets = {}
connection_lock = threading.Lock()

def handle_client(client_socket, client_addr):
    print(f"New connection from {client_addr}")
    client_type = None
    try:
        # First, receive the client's type identifier
        client_type = client_socket.recv(1024).decode('utf-8').strip()
        with connection_lock:
            client_sockets[client_socket] = client_type
        print(f"Registered as {client_type} from {client_addr}")

        while True:
            # Incoming message from client
            data = client_socket.recv(1024).decode('utf-8')
            if not data:
                break  # Client disconnected

            if client_type == 'controller':
                # Broadcast controller's task to all worker clients
                print(f"Broadcasting task: {data}")
                with connection_lock:
                    for sock, typ in client_sockets.items():
                        if typ == 'worker' and sock != client_socket:
                            try:
                                sock.sendall(data.encode('utf-8'))
                            except Exception as e:
                                print(f"Failed to send to worker: {e}")
            elif client_type == 'worker':
                # Forward worker's result to the controller
                print(f"Forwarding worker result: {data}")
                with connection_lock:
                    for sock, typ in client_sockets.items():
                        if typ == 'controller':
                            try:
                                sock.sendall(data.encode('utf-8'))
                            except Exception as e:
                                print(f"Failed to send to controller: {e}")
    except Exception as e:
        print(f"Error handling {client_addr}: {e}")
    finally:
        # Clean up disconnected client
        with connection_lock:
            if client_socket in client_sockets:
                del client_sockets[client_socket]
        client_socket.close()
        print(f"Connection closed with {client_addr}")

def start_server(host='localhost', port=12345):
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server_socket.bind((host, port))
    server_socket.listen(5)
    print(f"Server running on {host}:{port}")

    while True:
        client_sock, client_addr = server_socket.accept()
        # Spin up a thread to handle this client
        client_thread = threading.Thread(target=handle_client, args=(client_sock, client_addr))
        client_thread.daemon = True
        client_thread.start()

if __name__ == "__main__":
    start_server()

2. Controller Client Implementation

The controller connects to the server, sends task messages, and listens for worker results in a separate thread (so input and receiving don't block each other).

import socket
import threading

def listen_for_results(client_socket):
    while True:
        try:
            result = client_socket.recv(1024).decode('utf-8')
            if not result:
                print("Server disconnected")
                break
            print(f"Received worker result: {result}")
        except Exception as e:
            print(f"Error receiving results: {e}")
            break

def start_controller(host='localhost', port=12345):
    controller_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    controller_sock.connect((host, port))
    # Identify ourselves to the server
    controller_sock.sendall('controller'.encode('utf-8'))
    print("Connected to server as controller")

    # Start a thread to listen for worker results
    result_thread = threading.Thread(target=listen_for_results, args=(controller_sock,))
    result_thread.daemon = True
    result_thread.start()

    # Send tasks to broadcast
    while True:
        task = input("Enter task to broadcast (type 'exit' to quit): ")
        if task.lower() == 'exit':
            break
        controller_sock.sendall(task.encode('utf-8'))

    controller_sock.close()
    print("Controller disconnected")

if __name__ == "__main__":
    start_controller()

3. Worker Client Implementation

Workers connect to the server, wait for task broadcasts, execute tasks, and send results back. The execute_task function is where you'll plug in your actual business logic.

import socket
import threading
import time

def execute_task(task):
    # Replace this with your actual task logic (data processing, etc.)
    print(f"Executing task: {task}")
    time.sleep(2)  # Simulate task processing time
    return f"Worker {socket.gethostname()} completed task: '{task}'"

def listen_for_tasks(client_socket):
    while True:
        try:
            task = client_socket.recv(1024).decode('utf-8')
            if not task:
                print("Server disconnected")
                break
            # Run task and send result back to server
            result = execute_task(task)
            client_socket.sendall(result.encode('utf-8'))
        except Exception as e:
            print(f"Error handling task: {e}")
            break

def start_worker(host='localhost', port=12345):
    worker_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    worker_sock.connect((host, port))
    # Identify ourselves to the server
    worker_sock.sendall('worker'.encode('utf-8'))
    print("Connected to server as worker")

    # Start thread to listen for tasks
    task_thread = threading.Thread(target=listen_for_tasks, args=(worker_sock,))
    task_thread.daemon = True
    task_thread.start()

    # Keep worker running (wait for interrupt to exit)
    try:
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        print("Worker shutting down...")
        worker_sock.close()

if __name__ == "__main__":
    start_worker()
Key Notes for Reliability
  • Thread Safety: The connection_lock ensures we don't have race conditions when modifying the client_sockets dictionary across threads.
  • Error Handling: All socket operations include exception handling to gracefully handle disconnects or network issues without crashing the server/clients.
  • Scalability: If you need to add more workers later, just run additional instances of the worker script—no server changes required.
  • Binary Data: If you're sending non-text data (like files or serialized objects), skip the decode/encode steps and work directly with bytes.

内容的提问来源于stack exchange,提问作者Sourabh Jaiswal

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 09:50:54