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

Python下Twitch聊天Socket连接的多线程关键词抓取异常排查

问题根源分析
  • 共享Socket冲突:socket.recv()是操作系统级别的互斥操作,多个线程同时读取同一个Socket时,每条消息只会被随机一个线程抢到读取,其余线程无法获取该条消息,自然会出现要么所有线程都没抢到带关键词的消息,要么刚好各自抢到的消息都命中关键词的异常现象。而且多线程读同一个Socket完全解决不了卡顿漏消息的问题,Socket卡顿是因为接收缓冲区无新数据,多线程并发读取也不会拿到额外内容。
  • 消息解析错误:你每次固定读取2048字节,很容易出现半条消息或者多条消息粘在一起的情况,直接按:分割会出现解析失败,漏过关键词。
  • 缺失依赖定义:代码中使用的bcolors类未定义,运行会直接报错。
修复方案

采用单线程读Socket+消息队列+多线程消费检测的架构,从根源解决Socket争抢问题:

  1. 单独启动一个线程持续读取Socket数据,处理粘包后把完整消息存入安全队列
  2. 多个工作线程只从队列取消息做关键词检测,不会争抢Socket资源
  3. 任意线程检测到关键词后触发全局停止事件,所有线程同步终止
修复后代码
import os
import time
import socket
import threading
import queue
from dotenv import load_dotenv

# 加载环境变量
load_dotenv()

# 定义输出颜色
class bcolors:
    OKGREEN = '\033[92m'
    OKCYAN = '\033[96m'
    ENDC = '\033[0m'

# Socket相关配置
server = "irc.chat.twitch.tv"
port = 6667
nickname = "frankied003"
token = os.getenv("TWITCH_TOKEN")
channel = "#xqcow"

# 创建Socket并连接
sock = socket.socket()
sock.connect((server, port))
sock.send(f"PASS {token}\n".encode("utf-8"))
sock.send(f"NICK {nickname}\n".encode("utf-8"))
sock.send(f"JOIN {channel}\n".encode("utf-8"))

# 消息缓冲区、队列、停止事件
buffer = ""
msg_queue = queue.Queue(maxsize=1000)
stop_event = threading.Event()

# 单线程读Socket,处理粘包放入队列
def read_socket_thread():
    global buffer
    while not stop_event.is_set():
        try:
            resp = sock.recv(2048).decode("utf-8")
        except:
            time.sleep(0.01)
            continue
        buffer += resp
        # 按行分割,只处理完整消息
        while "\n" in buffer:
            line, buffer = buffer.split("\n", 1)
            line = line.strip()
            if not line:
                continue
            if line.startswith("PING"):
                sock.send("PONG\n".encode("utf-8"))
                continue
            # 把完整消息放入队列
            try:
                msg_queue.put_nowait(line)
            except queue.Full:
                # 队列满了丢最旧的消息,避免阻塞
                msg_queue.get_nowait()
                msg_queue.put_nowait(line)

# 工作线程:从队列取消息检测关键词
def worker_thread(correct_answers):
    while not stop_event.is_set():
        try:
            line = msg_queue.get(timeout=0.1)
        except queue.Empty:
            continue
        try:
            username = line.split(":")[1].split("!")[0]
            message = line.split(":")[2]
            stripped_message = " ".join(message.split()).lower()
        except:
            # 解析失败跳过坏消息
            continue
        # 匹配到关键词
        if stripped_message in correct_answers:
            print(bcolors.OKGREEN + username + " - " + message + bcolors.ENDC)
            stop_event.set()
            return
        # 输出自己发的消息
        if username == nickname:
            print(bcolors.OKCYAN + username + " - " + message + bcolors.ENDC)

while True:
    consoleInput = input("输入问题的正确答案,多个答案用','分隔:")
    if consoleInput == "stop":
        break
    # 重置停止事件和队列
    stop_event.clear()
    while not msg_queue.empty():
        msg_queue.get_nowait()
    # 生成正确答案数组
    correctAnswers = consoleInput.split(",")
    correctAnswers = [answer.strip().lower() for answer in correctAnswers]
    # 启动Socket读线程
    read_thread = threading.Thread(target=read_socket_thread, daemon=True)
    read_thread.start()
    # 启动3个工作线程
    workers = []
    for _ in range(3):
        t = threading.Thread(target=worker_thread, args=(correctAnswers,), daemon=True)
        workers.append(t)
        t.start()
    # 等待停止事件触发
    stop_event.wait()
    # 等待所有线程结束
    for t in workers:
        t.join()
    read_thread.join()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 07:06:03