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

Angular中SSE实现异常:无法自动接收事件需刷新页面

问题:SSE无法自动更新测试客户端UI

我在测试应用中实现SSE(Server-Sent Events),服务端已配置完成,使用端点api/v1/sse/document。预期执行扫描操作后,扫描结果能同时显示在测试客户端和主应用,但实际测试客户端必须刷新页面才能看到更新。

控制台无报错,请求端点返回200状态码,但测试客户端无法自动接收事件,必须刷新才会显示变更。

相关代码

sse.service.ts

(注:该文件代码以图片形式提供,核心功能应为创建并返回SSE的Observable)

document-list.component.ts

public ngOnInit(): void {
    this.getDocuments();
    this.registerServerSentEvent();
}

public ngOnDestroy(): void {
    this.closeServerSentEvent();
}

/**
 * getDocuments: 调用API服务获取文档列表
 */
public getDocuments(): void {
    this.aService.getDocuments().subscribe((documents: Documents[]) => {
        this.documents = documents;
    });
}

/**
 * markDocumentAsProcessed: 调用API服务标记文档为已处理
 */
public markDocumentAsProcessed(document: Documents): void {
    this.aService.markDocumentProcessed(document.document.id).subscribe({
        next: () => {
            // 从列表中移除已处理文档
            this.documents = this.documents.filter((doc) => doc.document.id !== document.document.id);
            this.closeDialog();
        },
        error: (error) => {
            console.log("markDocumentProcessed Error:", error);
            // 此处处理错误
        },
    });
}

/**
 * showDetails: 调用API服务获取文档图片并在弹窗中展示
 */
public showDetails(document: Documents): void {
    this.aService.getDocumentImage(document.document.id).subscribe((image: Blob) => {
        const url = window.URL.createObjectURL(image);
        const safeUrl: SafeUrl = this.sanitizer.bypassSecurityTrustUrl(url);
        this.selectedDocument = {...document, imageDataURL: safeUrl};
        this.displayDialog = true;
    });
}

/**
 * closeDialog: 点击处理按钮后关闭弹窗
 */
private closeDialog(): void {
    this.displayDialog = false;
    this.selectedDocument = null;
}

private registerServerSentEvent(): void {
    const sseUrl = `${this.aService.config.url}api/v1/sse/document`;
    this.sseSubscription = this.sseService.getServerSentEvent(sseUrl).subscribe((event: MessageEvent) => {
        const documentEvent = JSON.parse(event.data);
        const eventType = documentEvent.type;
        const eventData = documentEvent.data;

        switch (eventType) {
            case "NewDocument": {
                // 处理新文档事件
                break;
            }

            case "ViewDocument": {
                // 处理查看文档事件
                break;
            }

            case "WatchlistStatusUpdate": {
                // 处理监控列表状态更新事件
                break;
            }

            case "DocumentProcessed": {
                // 处理文档已完成事件
                const processedDocumentId = eventData.documentId;
                this.updateProcessedDocument(processedDocumentId);
                break;
            }
        }
    });
}

private updateProcessedDocument(processedDocumentId: string): void {
    // 在文档列表中找到已处理的文档
    const processedDocumentIndex = this.documents.findIndex((doc) => doc.document.id === processedDocumentId);
    if (processedDocumentIndex !== -1) {
        // 从列表中移除已处理文档
        this.documents.splice(processedDocumentIndex, 1);
        // 根据需要更新其他UI逻辑或执行额外操作
    }
}

private closeServerSentEvent(): void {
    if (this.sseSubscription) {
        this.sseSubscription.unsubscribe();
        this.sseService.closeEventSource();
    }
}

a.service.ts

public getDocuments(): Observable<Documents[]> {
    let url = this.config.url;
    if (!url.endsWith("/")) {
        url += "/";
    }

    return this.http
        .get<Documents[]>(`${url}api/v1/document/`, {
            headers: {
                authorization: this.config.authorization,
            },
        })
        .pipe(
            map((data) => {
                return data;
            }),
            catchError((error) => {
                throw error;
            })
        );
}

/**
 * markDocumentProcessed: 调用API服务标记文档为已处理
 *
 */
public markDocumentProcessed(documentId: string): Observable<Documents[]> {
    let url = this.config.url;
    if (!url.endsWith("/")) {
        url += "/";
    }

    const requestBody = {
        DocumentId: documentId,
    };

    return this.http
        .post<Documents[]>(`${url}api/v1/document/processed`, requestBody, {
            headers: {
                authorization: this.config.authorization,
            },
        })
        .pipe(
            map((data) => {
                return data;
            }),
            catchError((error) => {
                throw error;
            })
        );
}

/**
 * getDocumentImage: 调用API服务获取文档图片
 *
 */
public getDocumentImage(documentId: string): Observable<Blob> {
    let url = this.config.url;
    if (!url.endsWith("/")) {
        url += "/";
    }

    return this.http.get(`${url}api/v1/document/${documentId}/image/Photo`, {
        responseType: "blob",
        headers: {
            authorization: this.config.authorization,
        },
    });
}

排查与解决建议

  1. 检查SSE服务端响应格式
    SSE要求响应头必须包含Content-Type: text/event-stream,且需禁用缓存(如设置Cache-Control: no-cache)。可在浏览器Network面板查看api/v1/sse/document请求的响应头是否符合要求。

  2. 确认sse.service.ts的实现逻辑
    确保getServerSentEvent方法正确创建EventSource实例,并处理连接状态。建议添加错误监听输出异常:

    getServerSentEvent(url: string): Observable<MessageEvent> {
        return new Observable((observer) => {
            const eventSource = new EventSource(url, {
                withCredentials: true // 需携带认证信息时开启
            });
    
            eventSource.onmessage = (event) => observer.next(event);
    
            eventSource.onerror = (error) => {
                console.error("SSE连接错误:", error);
                observer.error(error);
                eventSource.close();
            };
    
            eventSource.onopen = () => console.log("SSE连接已建立");
    
            return () => eventSource.close();
        });
    }
    
  3. 验证事件推送的触发时机
    确认扫描完成后,服务端确实推送了DocumentProcessed类型事件,且数据格式正确(包含type: "DocumentProcessed"和data.documentId)。可在Network面板查看SSE请求的响应内容,确认是否有对应事件。

  4. 检查Angular变更检测
    若事件已接收但UI未更新,可手动触发变更检测:

    import { ChangeDetectorRef } from '@angular/core';
    
    constructor(private cdr: ChangeDetectorRef) {}
    
    private updateProcessedDocument(processedDocumentId: string): void {
        const processedDocumentIndex = this.documents.findIndex((doc) => doc.document.id === processedDocumentId);
        if (processedDocumentIndex !== -1) {
            this.documents.splice(processedDocumentIndex, 1);
            this.cdr.detectChanges(); // 手动触发变更检测
        }
    }
    

    或改用BehaviorSubject存储列表,借助Observable自动触发变更:

    private documentsSubject = new BehaviorSubject<Documents[]>([]);
    public documents$ = this.documentsSubject.asObservable();
    
    public getDocuments(): void {
        this.aService.getDocuments().subscribe((documents) => {
            this.documentsSubject.next(documents);
        });
    }
    
    private updateProcessedDocument(processedDocumentId: string): void {
        const currentDocs = this.documentsSubject.value;
        const updatedDocs = currentDocs.filter(doc => doc.document.id !== processedDocumentId);
        this.documentsSubject.next(updatedDocs);
    }
    

    模板中使用async管道订阅:

    <div *ngFor="let doc of documents$ | async">...</div>
    
  5. 确认认证信息传递
    若SSE端点需认证,确保EventSource携带正确信息。EventSource默认不支持自定义请求头,若使用Token认证,可开启withCredentials: true(依赖Cookie),或用fetch模拟SSE实现自定义头。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 05:54:54