TT Lab
开始
学习 学习路径 课程

打造好服务的计算机科学 — 用测量重新学习教科书概念

找出、统计并堵住缝隙 — 线程、asyncio、SQLite

在 TT Lab 中继续学习

目标

用字节码确认 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": 금액}, …]}

步骤

  1. 在 /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)。
  2. 在 /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。
  3. 在同一个文件里编写 locked_increment(threads, loops)。做与第 2 步相同的读取 → time.sleep(0) → 写入,但用 threading.Lock 封住,使 lost 为 0。不要删掉让出(time.sleep(0))——评分器会统计它被调用的次数。
  4. 在同一个文件里编写 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 的对象)。
  5. 在同一个文件里编写 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。
  6. 在同一个文件里编写 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 返回的对象)。
  7. 在同一个文件里编写 sqlite_lost(path, rounds, mode)。按下面的“SQLite 规则”,在一个线程里交替使用两个连接。用 params.json 的 sqlite_rounds 运行三种模式,以 rounds、read_modify_write、atomic、deferred_txn 写入 /root/svccs/concurrency/sqlite.json。
  8. 在同一个文件里编写 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

参考

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,所以能用线程重叠起来。