基于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:
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.
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()
- Thread Safety: The
connection_lockensures we don't have race conditions when modifying theclient_socketsdictionary 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/encodesteps and work directly with bytes.
内容的提问来源于stack exchange,提问作者Sourabh Jaiswal

