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

WSO2 Micro Integrator处理text/event-stream响应的支持与方案问询

WSO2 Micro Integrator 处理 text/event-stream (SSE) 响应指南

1. 原生内置支持情况

WSO2 Micro Integrator(MI)没有原生内置对text/event-stream(SSE,Server-Sent Events)这类持续流式响应的完整支持。MI核心基于请求-响应模型设计,默认会等待完整响应体接收完成后再执行后续处理逻辑,无法直接适配SSE这种单向、持续推送的流式协议。

2. 内置中介器支持情况

由于MI无原生SSE支持,不存在专门用于处理此类流式响应的内置中介器。所有处理逻辑都需要通过自定义扩展或变通方案实现。

3. 变通方案与自定义实现

以下是几种可行的处理方案,核心思路是绕过MI默认的全量响应接收逻辑,直接处理流式数据:

方案一:自定义Class Mediator

编写Java类实现流式读取与处理SSE事件,是最灵活的方案。

步骤1:实现自定义中介器类

import org.apache.synapse.MessageContext;
import org.apache.synapse.mediators.AbstractMediator;
import java.io.BufferedReader;
import java.io.InputStreamReader;
import org.apache.http.HttpEntity;
import org.apache.http.HttpResponse;
import org.apache.synapse.transport.http.HttpCoreConstants;

public class SSEStreamHandler extends AbstractMediator {
    @Override
    public boolean mediate(MessageContext mc) {
        try {
            // 获取HTTP响应对象
            HttpResponse httpResponse = (HttpResponse) mc.getProperty(HttpCoreConstants.HTTP_RESPONSE);
            HttpEntity entity = httpResponse.getEntity();
            BufferedReader reader = new BufferedReader(new InputStreamReader(entity.getContent()));
            
            String line;
            StringBuilder eventBuffer = new StringBuilder();
            
            // 逐行读取SSE流
            while ((line = reader.readLine()) != null) {
                if (line.startsWith("data:")) {
                    // 收集事件数据(SSE事件以data:开头)
                    eventBuffer.append(line.substring(5).trim());
                } else if (line.isEmpty()) {
                    // 空行表示一个事件结束,触发处理逻辑
                    processSSEEvent(eventBuffer.toString(), mc);
                    eventBuffer.setLength(0);
                }
            }
            
            // 处理最后一个未以空行结尾的事件
            if (eventBuffer.length() > 0) {
                processSSEEvent(eventBuffer.toString(), mc);
            }
        } catch (Exception e) {
            handleException("Failed to process SSE stream", e, mc);
        }
        return true;
    }

    // 自定义事件处理逻辑
    private void processSSEEvent(String eventData, MessageContext mc) {
        // 示例:打印事件日志,可扩展为转发到下游、存储等逻辑
        log.info("Received SSE event: " + eventData);
        mc.setProperty("SSE_EVENT_DATA", eventData);
    }
}

步骤2:打包并部署类文件
将编译后的class文件打包为JAR,放到MI的repository/components/lib目录下,重启MI。

步骤3:在序列中配置Class Mediator

<sequence name="SSEProcessingSeq" xmlns="http://ws.apache.org/ns/synapse">
    <call>
        <endpoint>
            <address uri="https://your-sse-api.example.com/stream" format="rest"/>
        </endpoint>
    </call>
    <!-- 调用自定义中介器处理SSE流 -->
    <class name="com.yourdomain.mediators.SSEStreamHandler"/>
    <!-- 后续可添加其他中介器处理事件数据 -->
</sequence>

方案二:使用Script Mediator(Groovy/JavaScript)

无需编写Java类,直接用脚本实现流式处理,适合快速原型开发。

Groovy脚本示例

<sequence name="SSEScriptSeq" xmlns="http://ws.apache.org/ns/synapse">
    <call>
        <endpoint>
            <address uri="https://your-sse-api.example.com/stream" format="rest">
                <!-- 开启流式传输,避免MI缓存全量响应 -->
                <property name="http.streaming" value="true"/>
            </address>
        </endpoint>
    </call>
    <script language="groovy"><![CDATA[
        def httpResponse = context.getProperty(HttpCoreConstants.HTTP_RESPONSE)
        def entity = httpResponse.getEntity()
        def reader = new BufferedReader(new InputStreamReader(entity.getContent()))
        
        String line
        StringBuilder eventBuffer = new StringBuilder()
        
        while ((line = reader.readLine()) != null) {
            if (line.startsWith("data:")) {
                eventBuffer.append(line.substring(5).trim())
            } else if (line.isEmpty()) {
                log.info("Processed SSE event: " + eventBuffer.toString())
                // 自定义事件处理逻辑
                eventBuffer.setLength(0)
            }
        }
        
        if (eventBuffer.length() > 0) {
            log.info("Processed final SSE event: " + eventBuffer.toString())
        }
    ]]></script>
</sequence>

关键配置注意事项

  • 调整超时设置:SSE是长连接,需在端点配置中增加超时参数避免连接被提前关闭:
    <property name="http.connection.timeout" value="3600000"/> <!-- 1小时 -->
    <property name="http.socket.timeout" value="3600000"/>
    
  • 响应头设置:若需将SSE事件转发到下游客户端,需设置响应头Content-Type: text/event-stream:
    <property name="messageType" value="text/event-stream" scope="transport"/>
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 13:15:04