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
相关产品推荐
相关产品推荐

