全栈开发

drssionpage实战qwen自动化

本实战主要是在qwen官网实现提交问并获取回到结果保存到本地

项目结构:

config.py

from pathlib import Path

root_path = Path(__file__).parent


class AuditConfig:
    # 基础文件路径控制
    INPUT_DIR = root_path / "input"
    OUTPUT_DIR = root_path / "output"
    TEMPLATE_MD_PATH = root_path / "prompt/prompt_template.md"
    BROWSER_PROFILE_DIR = r"./browser_profile"

    # 流水线核心控制参数
    HEADLESS_MODE = False  # 是否开启无头模式
    MAX_RESPONSE_TIMEOUT = 480  # 模型生成最大容忍超时(秒)

    # DeepSeek 网页高精 XPATH 状态定位符
    # 核心状态图标 //button[contains(@class, 'stop-button')]
    SEND_ICON_XPATH = "xpath://button[contains(@class, 'send-button')]"  # 发送/流结束(上箭头)
    STOP_ICON_XPATH = "xpath://button[contains(@class, 'stop-button')]"  # 深度思考中/流流出(正方形)

    # 文本内容提取器多轨兜底策略
    LAST_MSG_CONTAINER_XPATCH = "xpath://div[@class='qwen-markdown']"

file_handler.py

import json
import os
import re
from config import AuditConfig


class DataFileHandler:

    @staticmethod
    def save_batch_log(output_dir: str, filename: str, content: str):
        """安全留存带有详细推理和文本日志的批次底稿"""
        os.makedirs(output_dir, exist_ok=True)
        target_path = os.path.join(output_dir, filename)
        with open(target_path, "w", encoding="utf-8") as f:
            f.write(content)

    @classmethod
    def extract_failed_qa(cls):
        cfg = AuditConfig()
        txt_files = cfg.OUTPUT_DIR.glob("*.txt")
        result = []

        for txt_file in txt_files:
            with open(str(txt_file), "r", encoding="utf-8") as f:
                lines = f.readlines()
                text = "".join(lines)
                text = re.sub(r"\s+", "", text)
                json_str = text.split("</thinking>")[-1]
                json_obj = json.loads(json_str)
                if json_obj:
                    for obj in json_obj:
                        if obj["fail_reasons"]:
                            obj["file_name"] = txt_file.name
                            result.append(obj)

        json.dump(result, open(cfg.OUTPUT_DIR / "failed_qa.json", "w", encoding="utf-8"), ensure_ascii=False)

browser_manager.py

import time

from DrissionPage.common import Keys

from config import AuditConfig
from DrissionPage import ChromiumPage, ChromiumOptions


class QwenBrowserManager:
    def __init__(self, config: AuditConfig):
        self.cfg = config
        self.page = self.__init__browser()

    def __init__browser(self):
        """内置初始化底层Chrome内核"""
        co = ChromiumOptions()
        co.set_user_data_path(self.cfg.BROWSER_PROFILE_DIR)
        if self.cfg.HEADLESS_MODE:
            co.headless()
        return ChromiumPage(co)

    def _ensure_environment_state(self):
        """
        确保大模型处于:思考 模式
        """
        time.sleep(1.2)
        print("【思考】模式校验...")
        qwen_think_btn = self.page.ele(
            "xpath://main//div[@class='message-input-container-area']//div[@class='message-input-right-button']//div[@class='qwen-thinking-selector']")
        if qwen_think_btn is not None and "思考" not in qwen_think_btn.text:
            qwen_think_btn.click()
            time.sleep(0.5)
            think_span = self.page.ele("xpath://span[text()='思考']")
            if think_span:
                print("选择【思考】模式")
                think_span.click()
                time.sleep(0.5)
            else:
                raise ValueError("未找到【思考】模式页面元素")
        else:
            print("当前已是【思考】模式")

    def _check_and_handle_login(self):
        """
        https://chat.qwen.ai/auth
        判断当前是否掉线或跳转到了登录页,如果是,强行通过input()熔断流水线,等待人工介入
        如果已登录访问登录页面会自动跳转到https://chat.qwen.ai/
        """
        is_triggered = False
        # 使用while循环形成死锁保护,防止用户未完成登录就误敲回车跳出
        while True:
            login_btn = self.page.ele("tag:button@@text():登录", timeout=1)
            if login_btn:
                is_triggered = True
                login_btn.click()
                print("\n" + "!" * 60)
                print("🚨 [SECURITY ALERT] 检测到当前处于【未登录】或【 Session 已过期】状态!")
                print(f"🔗 当前页面已强制跳转至: {self.page.url}")
                print("🛑 自动化流水线已安全挂起保护。")
                print("🛠️  请立即在浏览器弹窗中:手动完成登录 / 解决图形验证码。")
                print("!" * 60)
            else:
                break
            # 触发终端阻塞,等待人工确认
            input("\n🎯 确认登录成功并重新进入 Chat 主页后,请在此处按【回车键(Enter)】恢复流水线...")

            # 缓冲 2 秒给页面留出重定向和加载新 DOM 的时间,然后 while 会再次校验 URL
            time.sleep(2)
        if is_triggered:
            # 如果从登录页面跳转过来,需要检查当前是不是思考模式
            print("🔄 [BROWSER] 登录成功重返主页,正在为您重设环境参数...")
            self._ensure_environment_state()

    def open_new_session(self, batch_idx: int):
        """清空旧的Chat ID 引用 通过强行重定向开启隔离的上下文对话"""
        print(f"\n🔄 [BROWSER] 正在重置并建立第 {batch_idx} 阶段全新上下文对话...")
        self.page.get("https://chat.qwen.ai")
        time.sleep(2)

        # 1. 第一步先看需不需要登录
        self._check_and_handle_login()

        # 2. 第二步如果没有触发登录,需要检查当前思考模式是否开启
        self._ensure_environment_state()

    def drag_in_file(self, file_paths: list):
        self.page.actions.drag_in("xpath://textarea", files=file_paths)
        print("拖拽上传文件")

    def upload_file(self, file_paths: list):
        # 1. 点击上传按钮图标
        upload_icon = self.page.ele("xpath://div[@class='message-input-container-area']//span[@class='anticon']")
        if not upload_icon:
            print("未找到上传按钮图标")
            return

        upload_icon.click()
        print("已点击上传按钮,等待菜单出现")

        # 2. 等待菜单出现并点击"上传附件"
        upload_menu_item = self.page.ele(
            "xpath://li[contains(@class, 'ant-dropdown-menu-item')]//span[text()='上传附件']")
        if not upload_menu_item:
            print("未找到:上传附件菜单项")
            return

        print("已点击上传附件")

        upload_menu_item.click.to_upload(file_paths=file_paths)

    def send_message(self, message: str, time_sleep: float = 0.8):
        """精准控制输入框聚焦、清理键盘流发送"""
        input_box = self.page.ele("xpath://textarea")
        input_box.click()
        input_box.clear()
        input_box.input(message)
        time.sleep(time_sleep)
        input_box.input(Keys.ENTER)
        wait_time = 0
        while True:
            if self.page.ele(self.cfg.STOP_ICON_XPATH, timeout=3):
                break
            else:
                print(f"等待上传结束...{time_sleep}s 后重试")
                time.sleep(time_sleep)
                input_box.input(Keys.ENTER)
                wait_time += time_sleep
            if wait_time > 60:
                raise RuntimeError("发送等待时间超过60秒,重试当前流程...")

    def wait_for_ai_complete(self, max_timeout: int) -> float:
        """
        异步状态机,实时监听qwen状态控制区
        """
        time.sleep(1.5)  # 容忍网页框架完成 DOM 数变换

        # 💡 策略 2:提交问题后的黄金窗口,立刻判断是否因为高频或风控被踢到了登录页
        self._check_and_handle_login()

        # 1. 监控正方形图标(是否进入深度思考或生成状态)
        # stop_btn = self.page.ele("xpath://div[@class='message-input-container']//div[@class='message-input-right-button']//div[@class='message-input-right-button-send']/button")
        # print(stop_btn)
        if self.page.ele(self.cfg.STOP_ICON_XPATH, timeout=3):
            print(">>> 📡 状态切换:大模型已开始深度思考/流式吐字...")
        else:
            self._check_and_handle_login()
            print(">>> ⚠️ 未在短时间内捕获到正方形状态,判定网络微卡,跳过直接监控结束标志...")

        # 2. 轮询监控上箭头(判断生成完全结束)
        start_time = time.time()
        while True:
            # 💡 策略 3:在长达数分钟的 AI 流式吐字过程中,持续监控是否突发掉线或会话崩塌
            self._check_and_handle_login()

            if self.page.ele(self.cfg.SEND_ICON_XPATH, timeout=0.5):
                print(">>> 🎉 状态复原:捕获到【发送箭头】,数据流生成完毕。")
                break
            if time.time() - start_time > max_timeout:
                print(f"\n❌ [BROWSER] 达到长生成保护限制阈值({max_timeout}秒),强行斩断监控。")
                break
            time.sleep(0.5)
        return time.time() - start_time

    def fetch_last_response_text(self) -> str:
        msg_containers = self.page.eles(self.cfg.LAST_MSG_CONTAINER_XPATCH, timeout=3)
        if not msg_containers:
            return "⚠️ [警告] 状态未就绪:页面上未检测到任何 message 容器。"
        current_container = msg_containers[-1]
        resp = current_container.text.strip()
        return resp

main.py

import time

from browser_manager import QwenBrowserManager
from config import AuditConfig
from file_handler import DataFileHandler


def run_pipeline():
    print("=" * 60)
    print("🌟 Qwen 数据集质检全自动流水线START 🌟")
    print("=" * 60)

    # 1. 实例化依赖组件
    cfg = AuditConfig()
    browser = QwenBrowserManager(cfg)

    data_files = cfg.INPUT_DIR.glob("*.txt")

    batch_index = 1
    for file_path in data_files:
        retry_count = 0
        max_retries = 3

        while retry_count < max_retries:
            try:
                # 2. 初始化首次会话
                browser.open_new_session(batch_idx=batch_index)
                file_paths = [cfg.TEMPLATE_MD_PATH, file_path.__str__()]
                browser.drag_in_file(file_paths=file_paths)
                browser.send_message("请严格按照.md附件的模板,质检.txt中数据", 5)
                # 状态机接入监控
                cost_time = browser.wait_for_ai_complete(max_timeout=cfg.MAX_RESPONSE_TIMEOUT)
                print(f"✅ [PIPELINE] 本批次计算流结束,大模型实际耗时: {cost_time:.2f} 秒")
                latest_raw_answer = browser.fetch_last_response_text()
                DataFileHandler.save_batch_log(cfg.OUTPUT_DIR.__str__(), file_path.name, latest_raw_answer)
                break
            except Exception as e:
                retry_count += 1
                print(f"❌ [ERROR] 处理文件 {file_path.name} 时出错 (第{retry_count}次重试): {str(e)}")

        if retry_count >= max_retries:
            raise RuntimeError(f"重试3次仍未通过,请人工排查:{file_path}")
    batch_index += 1
    time.sleep(3)
    print("\n" + "=" * 50)
    print("汇总异常QA数据集")
    DataFileHandler.extract_failed_qa()
    print("\n" + "=" * 50)
    print("🎉 [SUCCESS] 全量原始数据集已100%完全无缝、完全自动化处理完毕!")
    print("=" * 50)


if __name__ == '__main__':
    run_pipeline()