找出、统计并堵住缝隙 — 线程、asyncio、SQLite
目标
用字节码确认 counter += 1 为什么不是原子的,并亲手做出并统计线程中丢失的更新、座位重复预订、asyncio 中 await 之间的超额取款、SQLite 两个连接的更新丢失,再把它们修好。把阻塞事件循环的调用的代价用“事件循环延迟”测出来,并测量在 GIL 之下线程和进程完成 CPU 任务的速度。
为什么重要
竞态条件产生于读取的值与写入的值之间的缝隙。在线程中是字节码指令之间,在 asyncio 中是 await,在数据库中是 SELECT 与 UPDATE 之间。只用锁罩住缝隙一部分的修复,运行时没有任何报错,只是数字不对,因此在评审中很难被发现。所以本实验的评分器不只看你写下的数字,而是把你的函数放到不在材料里的输入上重新运行,与参考实现对照。涉及时间的判定,由于 Pod 共享 CPU,只在较宽的范围内判断。
材料
位于 /opt/fixtures/svccs/concurrency/ 之下。只读取,不要修改。
params.json threads·loops(2·3단계) · sqlite_rounds(7단계) · lag_handlers·lag_block_s·lag_tick_s(6단계)
cpu_n·io_tasks·io_sleep_s(8단계)
bookings.json 좌석 예약 요청 목록 [{"user": "u001", "seat": "A01"}, …] — 같은 좌석 요청이 몰려 온다
accounts.json {"balances": {계좌: 잔액}, "withdrawals": [{"account": 계좌, "amount": 금액}, …]}
步骤
- 在
/root/svccs/concurrency/bytecode.json中写入python(例如像 "3.12" 这样的主、次版本)、ops、atomic。ops是用dis.get_instructions拆解def bump(): global counter; counter += 1得到的指令名称列表,去掉以RESUME和RETURN开头的那些;atomic表示这个自增在线程之间是否是原子的(true/false)。 - 在
/root/svccs/concurrency/conc.py中编写unsafe_increment(threads, loops)。threads个线程各自对共享值自增loops次,每次自增按 读取 →time.sleep(0)→ 写入读到的值 + 1 的顺序进行。返回{"expected": threads×loops, "actual": 최종값, "lost": expected − actual}(占位符为最终值)。把用params.json的 threads、loops 运行得到的结果,以threads、loops、expected、actual、lost写入/root/svccs/concurrency/lost.json。 - 在同一个文件里编写
locked_increment(threads, loops)。做与第 2 步相同的读取 →time.sleep(0)→ 写入,但用threading.Lock封住,使lost为 0。不要删掉让出(time.sleep(0))——评分器会统计它被调用的次数。 - 在同一个文件里编写
book(requests, safe)。每个请求由一个线程执行 确认座位是否空闲 →time.sleep(0)→ 记录预订。当safe为真时用锁封住。返回{"requests", "seats"(서로 다른 좌석 수), "confirmed"(예약 성공 응답 수), "double_booked"(confirmed − 성공 응답에 나온 서로 다른 좌석 수)}(seats 是互不相同的座位数,confirmed 是预订成功的响应数,double_booked 是 confirmed 减去成功响应中出现的互不相同的座位数)。用bookings.json运行两次(不加锁、加锁),以requests、seats、unsafe、safe写入/root/svccs/concurrency/seats.json(unsafe、safe各是带有confirmed、double_booked的对象)。 - 在同一个文件里编写
run_withdrawals(balances, withdrawals, mode)。每笔取款用asyncio.gather按请求列表的顺序启动一个协程,每个协程执行 确认余额 ≥ 金额 →await asyncio.sleep(0)→ 扣减(余额不足则拒绝)。当mode为"locked"时,按账户用asyncio.Lock封住。返回{"approved", "rejected", "overdrawn_accounts"(최종 잔액이 음수인 계좌 수), "final"(계좌별 최종 잔액)}(overdrawn_accounts 是最终余额为负的账户数,final 是各账户的最终余额)。用accounts.json运行两种模式,以{"unsafe": …, "locked": …}写入/root/svccs/concurrency/overdraft.json。 - 在同一个文件里编写
async def handler(offload, block_s=0.2)和measure_lag(offload, handlers=3, tick=0.01)。handler 在offload为假时原样调用time.sleep(block_s),为真时通过asyncio.to_thread转移过去调用。measure_lag 按下面的“事件循环延迟规则”来测量。把两种情况以handlers、blocking、offloaded写入/root/svccs/concurrency/lag.json(blocking、offloaded就是 measure_lag 返回的对象)。 - 在同一个文件里编写
sqlite_lost(path, rounds, mode)。按下面的“SQLite 规则”,在一个线程里交替使用两个连接。用params.json的 sqlite_rounds 运行三种模式,以rounds、read_modify_write、atomic、deferred_txn写入/root/svccs/concurrency/sqlite.json。 - 在同一个文件里编写
cpu_bound(n)(从 0 到 n−1 的i*i % 7之和)、gil_compare(n)、io_compare(tasks, sleep_s),并按下面的“GIL 规则”写出/root/svccs/concurrency/gil.json。
事件循环延迟规则
ticker 작업: 반복해서 due = loop.time() + tick → await asyncio.sleep(tick) → loop.time() − due 를 기록
ticker 를 띄우고 tick×3 만큼 기다린 뒤, handler(offload) 를 handlers 개 gather 한다(걸린 시간 = elapsed_s)
gather 가 끝나면 ticker 를 멈추고 기다린다
돌려줄 것: max_lag_ms(기록 최댓값 × 1000, 소수 첫째 자리) · ticks(기록 개수) · elapsed_s(소수 셋째 자리)
SQLite 规则
path 에 새 DB: CREATE TABLE counter (id INTEGER PRIMARY KEY, n INTEGER NOT NULL), 행 (1, 0)
연결 a, b 는 sqlite3.connect(path, timeout=0.05). 라운드마다
read_modify_write a 가 n 을 읽고, b 가 n 을 읽고, a 가 (읽은 값+1) 을 쓰고 commit, b 도 같게
atomic a 가 UPDATE counter SET n = n + 1 후 commit, b 도 같게
deferred_txn a·b 의 isolation_level = None. a: BEGIN·읽기, b: BEGIN·읽기, a: 읽은 값+1 쓰기,
b: 읽은 값+1 쓰기 후 COMMIT — OperationalError 면 busy +1 하고 ROLLBACK, 끝으로 a: COMMIT
돌려줄 것: final(최종 n) · expected(2×rounds) · busy · lost(expected − final − busy)
GIL 规则
gil_compare(n): cpu_bound(n) 을 한 번 돌려 데운 뒤, cpu_bound(n) 두 번을 순서대로 / 스레드 2개로 /
ProcessPoolExecutor(max_workers=2) 로 돌린다. 셋 다 벽시계(perf_counter)와 CPU 시간을 함께 잰다
→ cpu_seq_s · cpu_threads_s · cpu_procs_s (벽시계 초)
→ seq_cpu_s · threads_cpu_s (time.process_time 차이 — 이 프로세스 모든 스레드의 CPU 시간)
→ procs_cpu_s (os.times() 의 children_user + children_system 차이 — 끝난 자식 프로세스의 CPU 시간)
(여섯 값 모두 소수 셋째 자리)
io_compare(tasks, sleep_s): time.sleep(sleep_s) tasks 번을 순서대로 / 스레드 tasks 개로 → io_seq_s · io_threads_s
gil.json: gil_disabled_build(sysconfig.get_config_var("Py_GIL_DISABLED") 를 bool 로), 위 여덟 값,
thread_speedup = cpu_seq_s ÷ cpu_threads_s, proc_speedup = cpu_seq_s ÷ cpu_procs_s,
io_thread_speedup = io_seq_s ÷ io_threads_s,
thread_cores = threads_cpu_s ÷ cpu_threads_s, proc_cores = procs_cpu_s ÷ cpu_procs_s
(다섯 비율 모두 소수 둘째 자리 — cores 는 '그동안 평균 몇 개의 코어가 일했나' 입니다)
lost_updates_with_gil = lost.json 의 lost
参考
- 评分器会导入
conc.py。生成结果 JSON 的代码,请放在if __name__ == "__main__":之下或单独的脚本里。 - 线程竞争的丢失数,由操作系统决定调度顺序,所以每次都不同。评分器检查的是:修复过的一侧是否恰好为 0,没有修复的一侧是否大于 0。asyncio 和 SQLite 一侧的数字每次都相同。
- 常见错误:只把写入用锁罩住;只把检查用锁罩住,而把操作放在外面;在协程里持有 threading.Lock 的同时 await(事件循环会停住);写了改用 to_thread,却原样保留 time.sleep。
- 产出会在会话结束后消失。需要的话请另行保存。
counter += 1 是几条指令
用 dis 拆解 counter += 1,把指令名称列表和原子性判断以 python、ops、atomic 的形式写入 /root/svccs/concurrency/bytecode.json。
dis.get_instructions(函数) 会逐条返回指令,每条指令的 opname 就是名称。如果读取(LOAD_)和写入(STORE_)是不同的指令,再结合“GIL 只保证单条指令”这一事实,想一想这说明了什么。
用线程丢失更新
在 conc.py 中编写 unsafe_increment(threads, loops),用 params.json 的 threads、loops 运行,把 threads、loops、expected、actual、lost 写入 /root/svccs/concurrency/lost.json。评分器也会用别的线程数、循环次数来调用这个函数。
把共享值放在全局变量或一个 dict 中,每个线程重复执行 v = 值 → time.sleep(0) → 值 = v + 1。sleep(0) 会放下 GIL,给其他线程制造读到同一个值的缝隙。expected 是 threads × loops。
用锁封住整个缝隙
在 conc.py 中编写 locked_increment(threads, loops)。保持读取 → time.sleep(0) → 写入不变,用 threading.Lock 使 lost 为 0。评分器会用多种线程数来调用,也会统计 time.sleep 被调用的次数。
锁该罩住什么,就是这一步的全部。如果只罩住写入那一行,在等锁期间,已经读到的值就已经过时,会依次覆盖。请把从读取到写入看作一个整体。
座位重复预订——先检查后操作
在 conc.py 中编写 book(requests, safe),用 bookings.json 分别不加锁、加锁运行,把 requests、seats、unsafe、safe 写入 /root/svccs/concurrency/seats.json。评分器也会用别的请求列表来调用。
如果在“是否空闲?”与“预订”之间,另一个线程确认了同一个座位,两个请求都会收到成功响应。只把检查放进锁里、把记录放在外面,缝隙依然存在。把座位收集到成功响应列表里,就能数出 double_booked(list.append 是 FAQ 写明为原子的操作)。
await 之间的超额取款
在 conc.py 中编写 run_withdrawals(balances, withdrawals, mode),用 accounts.json 运行 unsafe、locked,写入 /root/svccs/concurrency/overdraft.json。评分器会用别的账户和取款列表调用两种模式,并与参考实现精确对照。
asyncio 只有一个线程,但在 await 处会把轮次交给别的协程。如果确认与扣减之间有 await,在这期间同一账户的另一笔取款就会看到同样的余额。锁用 asyncio.Lock 配合 async with,每个账户一把,把从确认到扣减罩住。如果在协程里持有 threading.Lock 的同时 await,事件循环就会停住。
一个阻塞调用让所有人停住
在 conc.py 中编写 handler(offload, block_s=0.2) 和 measure_lag(offload, handlers=3, tick=0.01),把两种情况以 handlers、blocking、offloaded 写入 /root/svccs/concurrency/lag.json。评分器会用自己的测量器重新测量你的 handler。
time.sleep 运行期间,事件循环线程无法唤醒任何其他任务。把三个 handler 一起 gather,三者会在同一轮事件循环中接连阻塞,延迟因此叠加。转移到 to_thread 的情形,活本身还是要干的,所以 elapsed_s 不会变成 0。
SQLite 两个连接——静默丢失与明确失败
在 conc.py 中按 SQLite 规则编写 sqlite_lost(path, rounds, mode),用 sqlite_rounds 运行三种模式,把 rounds、read_modify_write、atomic、deferred_txn 写入 /root/svccs/concurrency/sqlite.json。评分器会用别的轮数调用三种模式。
Python 的 sqlite3 在默认设置下只在 UPDATE 之前开启事务——用 SELECT 读到的值得不到保护。UPDATE ... SET n = n + 1 把读取和写入放在一条语句里。如果手动 BEGIN 的两个连接都读完之后再想写入,其中一个无法升级为写入,会收到“database is locked”。
GIL——既不变快,也不安全
在 conc.py 中编写 cpu_bound、gil_compare、io_compare,并按 GIL 规则写出 /root/svccs/concurrency/gil.json。lost_updates_with_gil 是第 2 步 lost.json 中的 lost。评分器会直接调用 io_compare,并对照 cpu_bound 的值。
只用 CPU 的函数几乎不会放下 GIL,所以两个线程是轮流运行的。进程有各自独立的解释器,所以 GIL 也是各自独立的。墙钟时间在隔壁 Pod 繁忙时会波动,但“CPU 时间 ÷ 墙钟时间”是这段时间里平均有几个核心在工作,所以 GIL 的痕迹更清晰。time.sleep 在等待期间会放下 GIL,所以能用线程重叠起来。