You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用信号量同步两个发送线程与接收线程?

解决交替发送+接收的同步问题:信号量的正确用法

首先,咱们先拆解你的问题核心:需要严格遵循 S1 → S2 → Rec → S1 → S2 → Rec... 的循环,每一轮必须等两个发送线程都完成一次,接收线程才执行,然后再开启下一轮。你的代码已经实现了S1和S2的交替,但没把接收线程的执行和下一轮发送的启动绑定起来,导致信号量累积,流程混乱。

你的代码问题分析

原来的semRecNotify每次S2完成就释放,但没有机制让发送线程等待接收完成——S1和S2会自顾自地循环发送,导致semRecNotify的计数不断增加。等接收线程开始执行时,会一次性处理多轮累积的信号,直接破坏了每轮一次接收的规则。

修正后的代码实现

我们需要构建一个闭环同步链,让每一轮的结束(Rec执行完成)成为下一轮的开始(S1启动)的条件。这里用Python的threading模块给出示例:

import threading
import time

# 控制每一轮发送的启动:初始允许第一轮S1启动
sem_start_round = threading.Semaphore(value=1)
# S2等待S1完成的信号
sem_s2_ready = threading.Semaphore(value=0)
# Rec等待S2完成的信号
sem_rec_ready = threading.Semaphore(value=0)

def sender_one():
    while True:
        sem_start_round.acquire()  # 等待上一轮接收完成,或第一轮直接启动
        print("Sender One")
        sem_s2_ready.release()  # 通知S2可以发送
        time.sleep(0.1)  # 模拟发送耗时,避免输出混在一起

def sender_two():
    while True:
        sem_s2_ready.acquire()  # 等待S1完成
        print("Sender Two")
        sem_rec_ready.release()  # 通知Rec可以执行
        time.sleep(0.1)

def receiver():
    while True:
        sem_rec_ready.acquire()  # 等待S2完成
        print("Receiver")
        sem_start_round.release()  # 释放信号,允许下一轮S1启动
        time.sleep(0.1)

# 启动线程
t1 = threading.Thread(target=sender_one)
t2 = threading.Thread(target=sender_two)
t3 = threading.Thread(target=receiver)

t1.start()
t2.start()
t3.start()

t1.join()
t2.join()
t3.join()

流程说明

这个代码形成了完美的同步闭环:

  1. 初始时sem_start_round值为1,S1直接获取信号量开始发送;
  2. S1完成后释放sem_s2_ready,触发S2执行;
  3. S2完成后释放sem_rec_ready,触发Rec执行;
  4. Rec完成后释放sem_start_round,让下一轮的S1启动。
    每一步都严格依赖前一步的完成,完全符合你要的S1→S2→Rec循环规则。

通用的同步问题解决思路(数学化/结构化方法)

这类问题属于顺序依赖型同步,可以用以下通用步骤解决:

  1. 明确执行序列:把流程拆解成清晰的阶段,比如你的场景是[S1, S2] → Rec循环,每个阶段必须按顺序完成;
  2. 构建依赖链:用同步原语(信号量、屏障、条件变量等)给每个步骤添加"启动条件"——只有前序步骤完成,当前步骤才能启动;
  3. 避免信号累积:确保每个同步原语的release和acquire是一一对应的,每轮中每个信号只会被触发一次,防止出现"多轮信号堆在一起"的情况;
  4. 可选:用屏障简化多线程等待:如果需要多个线程同时完成后再执行下一步(比如S1和S2都完成后Rec执行),可以用Barrier(屏障)来替代多个信号量。比如设置一个Barrier(2),让S1和S2都到达屏障后,再触发Rec执行。

举个用Barrier的简化示例:

import threading
import time

send_barrier = threading.Barrier(2)  # 等待两个发送线程都完成
rec_sem = threading.Semaphore(value=0)
next_round_sem = threading.Semaphore(value=1)

def sender_one():
    while True:
        next_round_sem.acquire()
        print("Sender One")
        send_barrier.wait()  # 等待S2完成
        rec_sem.acquire()  # 等待接收完成
        next_round_sem.release()

def sender_two():
    while True:
        print("Sender Two")
        send_barrier.wait()  # 等待S1完成
        rec_sem.release()  # 触发Rec执行

def receiver():
    while True:
        rec_sem.acquire()
        print("Receiver")
        rec_sem.release()  # 通知发送线程开启下一轮

核心总结

解决同步问题的关键是清晰定义每个步骤的依赖关系,用同步原语把这些依赖落地,确保流程不会出现"超前执行"或"信号堆积"的情况。测试时可以加入模拟延迟,更容易看出同步是否正确。

内容的提问来源于stack exchange,提问作者SimpleThings

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 07:06:57