引言

在网络编程中,Socket是实现网络通信的基础。本文将详细介绍如何使用Python的Socket模块实现多个客户端通过服务端进行实时通信的系统。这种架构类似于聊天室系统,服务端作为消息中转站,负责转发所有客户端之间的消息。

系统架构概述

本系统采用客户端-服务器(C/S)架构:

  • 服务端:运行在指定端口,监听客户端连接,负责消息转发

  • 客户端:连接到服务端,可以发送和接收消息

  • 通信协议:使用JSON格式进行数据交换

核心技术要点

1. Socket编程基础

python

# 创建TCP Socket
client_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
serve_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)

  AF_INET:IPv4地址族

  SOCK_STREAM:TCP协议,提供可靠的双向字节流

2. 多线程处理

为了实现同时发送和接收消息,我们使用多线程技术:

# 客户端启动接收线程
recv_thread = threading.Thread(target=receive_message, daemon=True)
recv_thread.start()

3. JSON数据序列化

使用JSON格式进行数据传输:

  json.dumps():将Python对象转换为JSON字符串

  json.loads():将JSON字符串转换为Python对象

顺序图解

Serve端代码及注释

import socket
import threading
import json


def handle_client(conn_socket, client_addr, conn_socket_dict):
    try:
        client_ip, client_port = client_addr
        print(f"客户端 {client_ip}:{client_port} 已连接")

        # 发送欢迎消息(纯JSON格式)
        welcome_data = {
            "connection": list(conn_socket_dict.keys()), #使用list形式是因为 JSON 只能序列化基本的数据类型(字符串、数字、列表、字典等),不能序列化 Python 的特殊对象。
            "message": "连接服务器成功",
            "your_port": client_port  # 告诉客户端他们的端口号
        }
        #当客户端第一次连接上服务器时,服务器发来消息,告诉在线可联系人的信息
        conn_socket.send(json.dumps(welcome_data).encode('utf-8'))

        while True:
            #从客户端socket接收最多1024字节的数据,解码为‘utf-8'协议
            recv_data = conn_socket.recv(1024).decode('utf-8')

            #如果为空则说明断联
            if not recv_data:
                break

            try:
                #json.loads() 将 JSON 格式的字符串转换为 Python 对象
                data = json.loads(recv_data)
                print(data)
                to_client_port = data['to']  # 获取目标端口号
                data['from'] = str(client_port)  # 告知对方 发送方端口号

                # 查找目标客户端
                target_found = False
                #遍历连接字典
                for addr, sock in conn_socket_dict.items():
                    #addr[1]指元组里存储的第二个值,即端口号
                    if addr[1] == to_client_port:  # 通过端口号匹配
                        sock.send(json.dumps(data).encode('utf-8')) #json.dumps()将 Python 对象转换为 JSON 格式的字符串   发送给目标端口号
                        target_found = True
                        break

                if not target_found:
                    error_msg = {
                        "type": "error",
                        "message": f"用户{to_client_port}已经离线"
                    }
                    conn_socket.send(json.dumps(error_msg).encode('utf-8'))

            except json.JSONDecodeError:
                error_msg = {
                    "type": "error",
                    "message": "无效的消息格式"
                }
                conn_socket.send(json.dumps(error_msg).encode('utf-8'))

    except Exception as e:
        print(f"处理客户端 {client_addr} 时出错: {str(e)}")
    finally:
        # 从字典中移除断开连接的客户端
        if client_addr in conn_socket_dict:
            del conn_socket_dict[client_addr]
        conn_socket.close()
        print(f"客户端 {client_addr} 已断开连接")


if __name__ == '__main__':
    # 创建一个socket对象,AF_INET指的是intent间的通讯,SOCK_STREAM指的是使用的tcp通讯协议
    serve_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)

    # 实现端口复用
    serve_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, True)
    #创建连接字典,保存对应的地址和socket连接
    conn_socket_dict = {}

    #绑定本机的端口
    ip = '127.0.0.1'
    port = 8080
    serve_socket.bind((ip, port))
    #设置监听数为5
    serve_socket.listen(5)
    print(f"服务器正在监听 {ip}:{port}")

    try:
        #使用while True 进行长连接
        while True:
            #accept() -> (socket object, address info) 以元组的形式返回
            conn_socket, client_addr = serve_socket.accept()
            print(f"新客户端连接: {client_addr}")
            #存入连接字典中
            conn_socket_dict[client_addr] = conn_socket
            #进行多线程编程,在发消息的同时也能收到消息
            thread = threading.Thread(
                target=handle_client,
                args=(conn_socket, client_addr, conn_socket_dict),
                daemon=True
            )
            thread.start()

    except KeyboardInterrupt:
        print("正在关闭服务...")
    finally:
        serve_socket.close()
        print("服务已退出")

Client端代码及注释:

import socket
import json
import threading
#创建socket对象
client_socket= socket.socket(socket.AF_INET,socket.SOCK_STREAM)

#服务器的ip和端口
serve_ip = "127.0.0.1"
serve_port =8080

#连接服务器
def connect(serve_ip, serve_port):
    try:
        client_socket.connect((serve_ip, serve_port))
        #接收第一次连接服务器发来的消息
        welcome_msg = client_socket.recv(1024).decode('utf-8')
        print('目前有的端口号',json.loads(welcome_msg)["connection"],json.loads(welcome_msg)["message"],'   你的端口号是: ',json.loads(welcome_msg)["your_port"])
        return True  # 添加返回值
    except Exception as e:
        print(f'服务器链接失败: {str(e)}')
        return False  # 添加返回值



def send_message():
    try:
        while True:
            to = input('请问你想发消息给谁:')
            msg = input('请输入:')

            if msg.lower() == 'exit':  #退出连接机制
                break
            #消息格式
            data = {
                'to': int(to),
                'msg': msg
                    }
            json_data = json.dumps(data)
            client_socket.send(json_data.encode('utf-8'))

    except (KeyboardInterrupt, Exception) as e:
        print(f"\n发生错误: {str(e)}")
    finally:
        client_socket.close()
        print("客户端已关闭")





def receive_message():
    while True:
        try:
            recv_msg=client_socket.recv(1024).decode('utf-8')
            json_msg=json.loads(recv_msg)
            if not json_msg:
                break
            print('你收到来自'+json_msg['from']+'的消息'+json_msg['msg'])
        except Exception as e:
            print(str(e))


def start():
    if connect(serve_ip, serve_port):
        # 先启动接收线程
        recv_thread = threading.Thread(target=receive_message, daemon=True)
        recv_thread.start()

        # 然后开始发送消息(这会阻塞主线程)
        send_message()
    else:
        print('连接失败,请重试')


if __name__=='__main__':
    start()





测试样例:

更多推荐