Python 网络与并发:Socket、线程、进程与协程

结合 Socket、SSH、线程锁、队列、进程通信与协程示例,整理一篇更适合 Python 3 的网络与并发入门文档。

  • 先理解网络通信与 Socket
  • 再理解线程与进程的区别
  • 最后理解协程和事件驱动为什么适合高并发 IO

一、Socket 是什么

可以把 Socket 理解为“网络通信的编程接口”。

  • 浏览器访问网站,本质上也是在做网络通信
  • HTTP 建立在 TCP 之上
  • SSH、Web 服务、聊天程序,本质上都离不开 Socket

一个最小的 TCP 服务端

import socket

def handle_request(client):
    # 接收客户端发来的数据,最多接收 1024 字节
    data = client.recv(1024)
    print("收到请求:", data.decode("utf-8", errors="ignore"))

    # 返回一个最简单的 HTTP 响应。
    body = "hello from python socket server"
    response = (
        "HTTP/1.1 200 OK\r\n"
        "Content-Type: text/plain; charset=utf-8\r\n"
        f"Content-Length: {len(body.encode('utf-8'))}\r\n"
        "\r\n"
        f"{body}"
    )
    # 将响应数据发送给客户端
    client.sendall(response.encode("utf-8"))

# 创建 TCP Socket 对象
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
# 绑定 IP 地址和端口号
server.bind(("127.0.0.1", 8000))
# 开始监听,最大挂起连接数为 5
server.listen(5)

while True:
    # 阻塞等待客户端连接,返回连接对象和客户端地址
    connection, address = server.accept()
    try:
        handle_request(connection)
    finally:
        # 每次请求处理完及时关闭连接。
        connection.close()

这个例子也解释了这句话:浏览器可以看成 Socket 客户端,服务器就是 Socket 服务端

二、一个简化版远程命令执行示例

旧附件里保存的是一个基于 Socket 的“远程执行命令”演示。它非常适合理解收发协议,但不适合直接上线使用。

服务端

import os
import socket

# 创建 Socket 对象
server = socket.socket()
# 绑定地址和端口
server.bind(("127.0.0.1", 9999))
# 开始监听
server.listen()

while True:
    # 接受客户端连接
    conn, addr = server.accept()
    print("new conn:", addr)

    while True:
        # 接收客户端发送的命令
        data = conn.recv(1024)
        if not data:
            break

        command = data.decode("utf-8")
        # 这里只是演示 Socket 通信,不要在生产环境直接执行外部输入命令。
        # 执行命令并获取输出结果
        command_result = os.popen(command).read() or "cmd has no output..."
        payload = command_result.encode("utf-8")

        # 先发送结果长度,方便客户端按长度接收。
        conn.send(str(len(payload)).encode("utf-8"))
        # 发送实际的命令执行结果
        conn.send(payload)

客户端

import socket

# 创建 Socket 对象
client = socket.socket()
# 连接到服务端
client.connect(("127.0.0.1", 9999))

while True:
    # 获取用户输入的命令
    cmd = input(">>: ").strip()
    if not cmd:
        continue

    # 发送命令给服务端
    client.send(cmd.encode("utf-8"))

    # 接收服务端发送的结果长度
    result_size = int(client.recv(1024).decode("utf-8"))
    received = b""
    # 循环接收数据,直到接收到的数据长度等于结果长度
    while len(received) < result_size:
        received += client.recv(1024)

    # 打印命令执行结果
    print(received.decode("utf-8"))

这个示例真正值得学什么

  • 客户端和服务端需要约定通信协议
  • 二进制长度信息常用于解决“消息边界”问题
  • 任何网络程序都要考虑编码、粘包、异常和权限控制

三、SSH 与 Paramiko 的关系

把 SSH、Socket 和 paramiko 放在一起理解最容易:

  • Socket 解决“怎么连”
  • SSH 解决“如何在安全协议下远程连接”
  • paramiko 是 Python 里常见的 SSH 客户端库

如果只是理解原理,手写 Socket 很有价值; 如果要真正做远程连接和自动化运维,更常见的是直接使用成熟库而不是自己造协议。

四、线程和进程到底怎么区分

这部分容易被口号化,这里直接给结论:

  • 进程是资源分配的基本单位
  • 线程是 CPU 调度的基本单位
  • 同一进程下的线程共享内存空间
  • 不同进程之间默认内存隔离

“线程快还是进程快” 这个问题不严谨

要看你处理的是哪类任务:

  • CPU 密集型:更适合多进程
  • IO 密集型:多线程或协程通常更合适

五、线程的基本用法

import threading
import time

# 定义线程要执行的任务函数
def worker(name):
    print(f"{name} started")
    # 模拟耗时操作
    time.sleep(1)
    print(f"{name} finished")

threads = []
# 创建并启动 3 个线程
for index in range(3):
    # target 指定任务函数,args 传递参数(必须是元组)
    thread = threading.Thread(target=worker, args=(f"t{index}",))
    thread.start()
    threads.append(thread)

for thread in threads:
    # join 用于等待所有子线程结束。主线程会阻塞在这里,直到子线程执行完毕
    thread.join()

print("all done")

守护线程

import threading
import time

# 定义后台任务函数
def background_job():
    while True:
        print("background running...")
        time.sleep(1)

# 创建守护线程,daemon=True 表示主线程退出时,该线程也会随之退出
thread = threading.Thread(target=background_job, daemon=True)
thread.start()
# 主线程等待 2 秒后退出
time.sleep(2)
print("main thread exit")

守护线程的含义是:主线程退出后,不再等待它完整执行结束

六、为什么线程里还需要锁

有一种说法是“Python 3 不用了,优化了”,这是不准确的。

即使在 Python 3 中,只要多个线程同时修改共享数据,依然可能出现竞争问题,锁仍然非常重要。

互斥锁

import threading

total = 0
# 创建互斥锁对象
lock = threading.Lock()

def increase():
    global total
    for _ in range(10000):
        # 对共享变量做修改时加锁,避免竞争条件。with 语句会自动获取和释放锁
        with lock:
            total += 1

# 创建 2 个线程同时执行 increase 函数
threads = [threading.Thread(target=increase) for _ in range(2)]
for thread in threads:
    thread.start()
for thread in threads:
    thread.join()

print(total)

递归锁 RLock

当同一线程需要在嵌套调用中重复获取同一把锁时,可以使用 RLock

import threading

# 创建递归锁对象,允许同一个线程多次获取锁
lock = threading.RLock()
num = 0
num2 = 0

def run1():
    global num
    # 获取锁
    with lock:
        num += 1
        return num

def run2():
    global num2
    # 获取锁
    with lock:
        num2 += 1
        return num2

def run3():
    # 获取锁,在 run3 内部调用 run1 和 run2 时,由于是同一个线程,可以再次获取锁而不会死锁
    with lock:
        first = run1()
        second = run2()
        print(first, second)

thread = threading.Thread(target=run3)
thread.start()
thread.join()

七、信号量与事件

信号量 Semaphore

信号量适合“限制同时并发数量”的场景。

import threading
import time

# 创建有界信号量,最多允许 5 个线程同时获取
semaphore = threading.BoundedSemaphore(5)

def run(index):
    # 获取信号量,如果当前已有 5 个线程获取了信号量,则阻塞等待
    with semaphore:
        # 同一时刻最多只有 5 个线程能进入这里。
        time.sleep(1)
        print(f"run the thread: {index}")

# 创建 10 个线程
threads = [threading.Thread(target=run, args=(i,)) for i in range(10)]
for thread in threads:
    thread.start()
for thread in threads:
    thread.join()

事件 Event

Event 适合做线程间的简单通知。

import threading
import time

# 创建事件对象,初始状态为 False
event = threading.Event()

def waiter():
    print("waiting...")
    # 阻塞等待事件状态变为 True
    event.wait()
    print("go")

thread = threading.Thread(target=waiter)
thread.start()
time.sleep(1)
event.set()  # 发出信号,将事件状态设置为 True,唤醒等待线程
thread.join()

八、线程队列与生产者消费者模型

以下几个队列类型很重要:

  • queue.Queue:先进先出
  • queue.LifoQueue:后进先出
  • queue.PriorityQueue:按优先级取出

生产者消费者示例

import queue
import threading
import time

# 创建一个最大容量为 10 的先进先出队列
task_queue = queue.Queue(maxsize=10)

def producer():
    count = 1
    while count <= 5:
        task = f"骨头{count}"
        # 将任务放入队列,如果队列已满则阻塞等待
        task_queue.put(task)
        print("生产了:", task)
        count += 1
        time.sleep(0.1)

def consumer(name):
    while True:
        # 从队列中取出任务,如果队列为空则阻塞等待
        item = task_queue.get()
        print(f"[{name}] 取到 {item} 并处理")
        time.sleep(0.5)
        # 标记任务处理完成
        task_queue.task_done()

# 创建并启动两个消费者线程,设置为守护线程
threading.Thread(target=consumer, args=("c1",), daemon=True).start()
threading.Thread(target=consumer, args=("c2",), daemon=True).start()
# 在主线程中运行生产者
producer()
# 阻塞等待队列中所有任务都被处理完成(即 task_done() 被调用的次数等于 put() 的次数)
task_queue.join()

生产者消费者模型的价值在于:

  • 解耦生产和消费速度
  • 平滑突发流量
  • 降低模块之间的直接依赖

补充:多线程下载脚本该怎么整理

临时目录里的旧脚本本质上是在做一件很典型的事:

  • 从 JSON 文件读取下载任务
  • 多线程并发拉取远程资源
  • 按文件名把内容落到本地目录

这个思路本身没有问题,但旧写法里混用了共享生成器、裸 except、手工 os.path 拼接等做法。整理到今天,更推荐写成下面这种结构更清晰的版本:

from __future__ import annotations

import json
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
import requests

HEADERS = {
    "User-Agent": (
        "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
        "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/126.0 Safari/537.36"
    )
}


def load_tasks(json_file: Path) -> list[tuple[str, str]]:
    data = json.loads(json_file.read_text(encoding="utf-8"))
    return [(item["Url"], item["FileName"]) for item in data]


def download_file(url: str, file_name: str, output_dir: Path) -> None:
    target = output_dir / file_name
    if target.exists():
        return

    with requests.get(url, stream=True, headers=HEADERS, timeout=10) as response:
        response.raise_for_status()
        if int(response.headers.get("content-length", "0")) == 0:
            raise ValueError(f"empty response body: {url}")

        with target.open("wb") as file:
            for chunk in response.iter_content(chunk_size=8192):
                if chunk:
                    file.write(chunk)


def main() -> None:
    output_dir = Path("output")
    output_dir.mkdir(exist_ok=True)
    tasks = load_tasks(Path("json1.json"))

    with ThreadPoolExecutor(max_workers=20) as executor:
        for url, file_name in tasks:
            executor.submit(download_file, url, file_name, output_dir)


if __name__ == "__main__":
    main()

这里保留了旧脚本真正值得学的点:

  • stream=True 适合下载大文件,避免一次性把内容全部读进内存
  • 先判断目标文件是否已存在,可以减少重复下载
  • 下载前确保输出目录存在

但也顺手替换掉了几个旧问题:

  • Path 代替手工拼路径
  • ThreadPoolExecutor 代替共享生成器加锁的调度方式
  • raise_for_status() 代替手工判断状态码区间
  • 避免使用裸 except 吞掉真实异常

九、多进程与进程间通信

multiprocessing.Queue

线程间用 queue.Queue,进程间更常用 multiprocessing.Queue

from multiprocessing import Process, Queue

def producer(queue_obj):
    # 子进程向队列中放入数据
    queue_obj.put("hello from child")

if __name__ == "__main__":
    # 在主进程中创建多进程安全的队列
    queue_obj = Queue()
    # 创建子进程,并将队列对象作为参数传递
    process = Process(target=producer, args=(queue_obj,))
    process.start()
    # 主进程从队列中获取数据
    print(queue_obj.get())
    process.join()

Pipe

Pipe 更像是双向管道,适合较明确的一对一通信。

from multiprocessing import Pipe, Process

def child(conn):
    # 子进程通过管道发送数据
    conn.send(["hello from child"])
    conn.close()

if __name__ == "__main__":
    # 创建管道,返回两个连接对象,分别代表管道的两端
    parent_conn, child_conn = Pipe()
    # 创建子进程,并将管道的一端传递给子进程
    process = Process(target=child, args=(child_conn,))
    process.start()
    # 主进程通过管道的另一端接收数据
    print(parent_conn.recv())
    process.join()

Manager

Manager 可以提供跨进程共享的代理对象。

from multiprocessing import Manager, Process

def worker(shared_dict, shared_list):
    # 子进程修改共享字典和列表
    shared_dict["name"] = "alex"
    shared_list.append("python")

if __name__ == "__main__":
    # 使用 Manager 创建一个管理器对象,用于管理共享数据
    with Manager() as manager:
        # 创建跨进程共享的字典和列表
        shared_dict = manager.dict()
        shared_list = manager.list()

        # 创建子进程,并将共享对象传递给子进程
        process = Process(target=worker, args=(shared_dict, shared_list))
        process.start()
        process.join()

        # 打印修改后的共享数据
        print(dict(shared_dict))
        print(list(shared_list))

十、协程与事件驱动

旧资料里把 yieldgreenletgeventselectors 放在一起,背后想表达的是:单线程也可以通过“遇到 IO 就切换”来提高并发效率

yield 理解最原始的协程思想

def consumer(name):
    while True:
        # yield 暂停执行,等待接收数据
        item = yield
        print(f"{name} is eating baozi {item}")

def producer():
    # 创建两个消费者协程
    c1 = consumer("c1")
    c2 = consumer("c2")
    # 先激活生成器,使其执行到 yield 处暂停
    next(c1)
    next(c2)

    for number in range(1, 4):
        # 通过 send() 方法向协程发送数据,并恢复其执行
        c1.send(number)
        c2.send(number)

producer()

gevent 的定位

gevent 是对协作式并发的一种封装,重点不是语法,而是理念:

  • 遇到 IO 时主动让出执行权
  • 单线程里也能同时处理很多连接
  • 更适合高并发 IO 场景,而不是纯计算

注意:现代 Python 项目更推荐使用 asyncio + async/awaitgevent 已逐渐边缘化。

十一、现代 asyncio 与 HTTP 客户端(推荐)

Python 3.7+ 之后,优先使用 asyncio.run() 管理事件循环。抓取类 IO 并发可结合 httpxaiohttp

import asyncio
import httpx

# 定义异步函数,使用 async 关键字
async def fetch(client: httpx.AsyncClient, url: str) -> tuple[str, int]:
    # 使用 await 关键字等待异步操作完成,期间会让出执行权
    resp = await client.get(url, timeout=10)
    return url, resp.status_code

async def main():
    urls = [f"https://httpbin.org/status/{code}" for code in (200, 301, 404, 500)]
    # 使用异步上下文管理器创建 HTTP 客户端
    async with httpx.AsyncClient(http2=True, headers={"User-Agent": "demo"}) as client:
        # asyncio.gather 并发执行多个异步任务
        results = await asyncio.gather(*[fetch(client, u) for u in urls], return_exceptions=True)
    for item in results:
        print(item)

# 使用 asyncio.run() 启动事件循环并运行主协程
asyncio.run(main())

要点:

  • asyncio.run() 而不是手工获取事件循环
  • IO 密集型优先选择异步客户端(httpx / aiohttp
  • 结合 asyncio.Semaphore 做并发限速

十二、如何选择线程、进程、协程

可以用这张速记表:

  • 大量网络请求、爬虫、Socket 连接:优先考虑协程或线程
  • CPU 计算密集任务:优先考虑多进程
  • 需要跨机器、跨语言传递任务:考虑消息队列
  • 只是本进程内解耦生产和消费:先用 queue.Queue

小结

学完这一篇,你至少应该建立下面这些判断:

  • Socket 是网络程序的基础抽象
  • 线程共享内存,进程默认隔离内存
  • 锁在 Python 3 里依然重要
  • queue.Queuemultiprocessing.Queue 不要混用概念
  • 协程更适合高并发 IO,而不是替代所有并发模型

理解这条主线以后,再去看爬虫、Web 服务器、消息队列和异步框架,会顺很多。

相关阅读

  • 《Python 消息队列与任务系统》:RabbitMQ、Redis、Celery