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

Python Flask SSE端点发送一次数据后连接断开的问题求助

问题描述

我们基于Python Flask实现了SSE(Server-Sent Events)端点/description/progress,但在React+TypeScript前端调用该端点时,服务器仅返回一次数据后就关闭连接。目前通过在响应中设置5秒重试来维持功能,但这属于暴力解决方案。先后尝试过用TypeScript原生逻辑处理,以及使用微软的fetch-event-source库(因原生EventSource无法设置请求头),但问题仍未解决。


初始TypeScript处理代码

export async function* descriptionProgressApi(fileName: string, idToken: string) {
    let response: Response;
    fileName = fileName.replace("Product_description_files/", "");
    try {
        response = await fetch(`/description/progress?` + new URLSearchParams({ file_name: fileName }), {
            method: "GET",
            headers: getHeaders(idToken)
        });
    } catch (networkError) {
        throw new Error(`Network error occurred: ${networkError}`);
    }
    if (!response.ok) {
        const errorDetails = await response.text();
        throw new Error(`Get progress of description failed: ${errorDetails}`);
    }
    if (!response.body) {
        throw new Error("Response body is null");
    }
    // Check if its processed, if it is then just yield and return
    const reader = response.body.getReader();
    const decoder = new TextDecoder();

    let done = false;
    let isFirstChunk = true;
    while (!done) {
        try {
            const { value, done: readerDone } = await reader.read();
            console.log(value);
            done = readerDone;
            if (value) {
                const chunk = decoder.decode(value, { stream: true });
                let cleanedChunk = chunk.replace(/^data: /, "");
                cleanedChunk = cleanedChunk.replace("retry: 5000", "");
                if (isFirstChunk) {
                    console.log("calling first chunk");
                    isFirstChunk = false;
                    try {
                        console.log(cleanedChunk);
                        const responseJSON = JSON.parse(cleanedChunk);
                        // Check if progress is 1 in the firstChunk
                        if (responseJSON.status === "Uploaded") {
                            yield responseJSON;
                            return;
                        } else if (responseJSON.status === "Not_initiated") {
                            yield responseJSON;
                            return;
                        }
                    } catch (error) {
                        console.error("Failed to parse initial JSON response:", error);
                        throw new Error("Initial JSON response parsing failed");
                    }
                }
                try {
                    console.log(cleanedChunk);
                    const json = JSON.parse(cleanedChunk);
                    if (json.status === "Started") {
                        yield json;
                    } else if (json.status === "Failed") {
                        yield json;
                        throw new Error("Error occured in uploading the file");
                    } else if (json.status === "Finished") {
                        console.log("returning finished");
                        yield json;
                        return;
                    }
                } catch (error) {
                    console.error("Failed to parse JSON:", error);
                }
            }
        } catch (error) {
            console.error("Stream reading error:", error);
            done = true; // Exit the loop if an error occurs
        }
    }
}

Flask后端代码

async def check_product_progress(product_file_id):
    try:
        session = SessionLocal()
        product_file = session.query(ProductFile).filter(ProductFile.ID == product_file_id).first()
        num_of_rows = product_file.number_of_rows
        processed = session.query(Product).filter(Product.processed == True, Product.product_file == product_file).count()
        print(num_of_rows, product_file, processed)

        if num_of_rows != 0:
            return {
                "file_name": product_file.Title,
                "progress": processed / num_of_rows,
                "number_of_rows": num_of_rows,
                "processed": processed,
                "status": product_file.status.name,
            }
        else:
            return {
                "file_name": product_file.Title,
                "progress": 0,
                "number_of_rows": num_of_rows,
                "processed": processed,
                "status": product_file.status.name,
            }
    except SQLAlchemyError as e:
        raise e
    finally:
        session.close()


async def description_progressbar(product_file_id, config):
    status = ProgressEnum(1)
    processing = status
    while status == processing:
        progress = await check_product_progress(product_file_id)
        event = ServerSentEvent(data=json.dumps(progress), retry=5000)
        yield event.encode()
        await asyncio.sleep(5)
        status = progress["status"]

    if progress["number_of_rows"] == progress["processed"] and status != ProgressEnum(3).name:
        await reassemble(product_file_id, config)
        progress = await check_product_progress(product_file_id)
        event = ServerSentEvent(json.dumps(progress))
        yield event.encode()


@bp.route("/description/progress", methods=["GET"])
@authenticated
async def description_progress(auth_claims: Dict[str, Any]):
    user_oid = auth_claims.get("oid", False)
    if not user_oid:
        return "Unauthorized", 401
    file_name = request.args.get("file_name", False)
    if file_name:
        session = SessionLocal()
        try:
            product_file = (
                session.query(ProductFile)
                .filter(ProductFile.Title == file_name, ProductFile.user_id == user_oid)
                .first()
            )
            if product_file:
                response = await make_response(
                    description_progressbar(product_file.ID, current_app.config), 
                    {
                    "Status": "Success", 
                    'Content-Type': 'text/event-stream',
                    'Cache-Control': 'no-cache',
                    'Transfer-Encoding': 'chunked',
                    'Connection': 'keep-alive'
                })
                response.timeout = None
                return response
            else:
                return "No such file", 400
        except SQLAlchemyError as e:
            session.rollback()
            raise e
        finally:
            session.close()
    else:
        return "No file provided", 400

SSE类代码

from dataclasses import dataclass


@dataclass
class ServerSentEvent:
    data: str
    event: str | None = None
    id: int | None = None
    retry: int | None = None

    def encode(self) -> bytes:
        message = f"data: {self.data}"
        if self.event is not None:
            message = f"{message}\nevent: {self.event}"
        if self.id is not None:
            message = f"{message}\nid: {self.id}"
        if self.retry is not None:
            message = f"{message}\nretry: {self.retry}"
        message = f"{message}\r\n\r\n"
        return message.encode("utf-8")

使用fetch-event-source的前端代码

async function handleProgressBla(fileName: string) {
        class RetriableError extends Error {}
        const idToken = await getToken(client);
        if (!idToken) {
            throw new Error("No authentication token available");
        }

        const headers = new Headers();
        headers.append("Authorization", `Bearer ${idToken}`);
        headers.append("Accept", `text/event-stream`);
        let data = 0;
        async function startEventSource() {
            try {
                await fetchEventSource("/description/progress?" + new URLSearchParams({ file_name: fileName }), {
                    headers: Object.fromEntries(headers),
                    async onopen(response) {
                        if (response.ok && response.status === 200) {
                            console.log("Connection opened!");
                        } else if (response.status >= 400 && response.status < 500 && response.status !== 429) {
                            console.log("Client-side error, won't retry:", response.statusText);
                        }
                    },
                    onmessage(event) {
                        const parsedData = JSON.parse(event.data);
                        console.log(parsedData);
                        if (parsedData.status === "Started") {
                            console.log("Started");
                            setProgressStatus(prevStatus =>
                                prevStatus.map(item => (item.name === "Product_description_files/" + fileName ? { ...item, status: "processing" } : item))
                            );
                        }

                        setProgress(parsedData.progress);
                    },
                    onclose() {
                        console.log("Connection closed by the server");
                        throw new RetriableError();
                       
                    },
                    onerror(err) {
                        console.log("There was an error from the server", err);
                        
                    }
                });
            } catch (error) {
                console.error("Failed to connect. Retrying...", error);
                // Optionally add a delay before reconnecting
                setTimeout(startEventSource, 2000); 
            }
        }

        startEventSource(); // Start the event source connection
    }

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 22:32:32