- 先理解网络通信与 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))
十、协程与事件驱动
旧资料里把 yield、greenlet、gevent、selectors 放在一起,背后想表达的是:单线程也可以通过“遇到 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/await,gevent已逐渐边缘化。
十一、现代 asyncio 与 HTTP 客户端(推荐)
Python 3.7+ 之后,优先使用 asyncio.run() 管理事件循环。抓取类 IO 并发可结合 httpx 或 aiohttp。
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.Queue和multiprocessing.Queue不要混用概念- 协程更适合高并发 IO,而不是替代所有并发模型
理解这条主线以后,再去看爬虫、Web 服务器、消息队列和异步框架,会顺很多。
相关阅读
- 《Python 消息队列与任务系统》:RabbitMQ、Redis、Celery