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

如何在Quarkus REST中访问HTTP Trailer以检测流传输完成?

在Quarkus REST中访问HTTP Trailer的实现方案

Quarkus REST基于RESTEasy Reactive(默认)或传统Jakarta REST,以下分服务端发送、客户端接收两种场景给出实现方式:

一、服务端发送HTTP Trailer

无论使用哪种模式,都需要先通过Trailer响应头声明将要发送的Trailer字段名,这是HTTP标准的强制要求。

1. Reactive模式(推荐)

使用@ServerTrailer注解声明Trailer,结合Multi实现流式响应:

import jakarta.ws.rs.GET;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.Produces;
import jakarta.ws.rs.core.MediaType;
import jakarta.ws.rs.core.Response;
import org.jboss.resteasy.reactive.RestStreamElementType;
import org.jboss.resteasy.reactive.ServerTrailer;
import io.smallrye.mutiny.Multi;
import java.util.Map;

@Path("/stream")
public class StreamingResource {

    @GET
    @Produces(MediaType.APPLICATION_OCTET_STREAM)
    @RestStreamElementType(MediaType.TEXT_PLAIN)
    @ServerTrailer(name = "Stream-Completed", description = "标记流是否成功完成")
    public Response streamData() {
        // 模拟流式数据输出
        var dataStream = Multi.createFrom().items("Chunk 1", "Chunk 2", "Chunk 3")
                .map(String::getBytes);

        // 定义Trailer内容,流结束时自动发送
        var trailers = Map.of("Stream-Completed", "true");

        return Response.ok(dataStream)
                .header("Trailer", "Stream-Completed") // 告知客户端存在该Trailer字段
                .trailers(() -> trailers)
                .build();
    }
}

2. 传统Jakarta REST模式

使用StreamingOutput结合HttpServletResponse手动添加Trailer:

import jakarta.ws.rs.GET;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.Produces;
import jakarta.ws.rs.core.MediaType;
import jakarta.ws.rs.core.Response;
import jakarta.ws.rs.core.StreamingOutput;
import jakarta.servlet.http.HttpServletResponse;
import jakarta.ws.rs.core.Context;
import java.io.IOException;
import java.io.OutputStream;

@Path("/stream-classic")
public class ClassicStreamingResource {

    @GET
    @Produces(MediaType.APPLICATION_OCTET_STREAM)
    public Response streamData(@Context HttpServletResponse response) {
        response.setHeader("Trailer", "Stream-Completed");

        StreamingOutput output = out -> {
            try {
                out.write("Chunk 1".getBytes());
                out.flush();
                out.write("Chunk 2".getBytes());
                out.flush();
                out.write("Chunk 3".getBytes());
                out.flush();
                // 流传输完成后添加Trailer
                response.addTrailerHeader("Stream-Completed", "true");
            } catch (Exception e) {
                response.addTrailerHeader("Stream-Completed", "false");
                throw new IOException(e);
            }
        };

        return Response.ok(output).build();
    }
}

二、客户端接收HTTP Trailer

1. Reactive REST Client

使用@ClientTrailer声明要接收的Trailer,并通过原方法名+Trailers的命名规则定义方法获取Trailer内容:

import jakarta.ws.rs.GET;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.Produces;
import jakarta.ws.rs.core.MediaType;
import org.eclipse.microprofile.rest.client.inject.RegisterRestClient;
import org.jboss.resteasy.reactive.ClientTrailer;
import org.jboss.resteasy.reactive.RestStreamElementType;
import io.smallrye.mutiny.Multi;
import java.util.Map;

@RegisterRestClient(baseUri = "http://localhost:8080")
public interface StreamingClient {

    @GET
    @Path("/stream")
    @Produces(MediaType.APPLICATION_OCTET_STREAM)
    @RestStreamElementType(MediaType.TEXT_PLAIN)
    @ClientTrailer(name = "Stream-Completed")
    Multi<String> getStream();

    // 严格遵循命名规则:原方法名 + Trailers
    @GET
    @Path("/stream")
    @Produces(MediaType.APPLICATION_OCTET_STREAM)
    Map<String, String> getStreamTrailers();
}

使用客户端时,在流完成后调用Trailer方法:

import jakarta.enterprise.context.ApplicationScoped;
import jakarta.inject.Inject;
import org.eclipse.microprofile.rest.client.inject.RestClient;
import io.smallrye.mutiny.Uni;

@ApplicationScoped
public class StreamConsumer {

    @Inject
    @RestClient
    StreamingClient client;

    public Uni<Void> consumeStream() {
        return client.getStream()
                .onItem().invoke(chunk -> System.out.println("收到数据块:" + chunk))
                .onCompletion().invoke(() -> {
                    var trailers = client.getStreamTrailers();
                    System.out.println("流传输完成,Trailer值:" + trailers.get("Stream-Completed"));
                })
                .collect().asList()
                .onItem().ignore().andContinueWithNull();
    }
}

2. 传统Jakarta REST Client

读取完流式响应后,直接从Response对象中获取Trailer:

import jakarta.ws.rs.client.Client;
import jakarta.ws.rs.client.ClientBuilder;
import jakarta.ws.rs.core.Response;
import java.io.InputStream;

public class ClassicClient {
    public static void main(String[] args) {
        Client client = ClientBuilder.newClient();
        try (Response response = client.target("http://localhost:8080/stream-classic").request().get()) {
            InputStream in = response.readEntity(InputStream.class);
            byte[] buffer = new byte[1024];
            int bytesRead;
            while ((bytesRead = in.read(buffer)) != -1) {
                System.out.println("收到数据块:" + new String(buffer, 0, bytesRead));
            }
            // 流读取完成后获取Trailer
            String completed = response.getHeaderString("Stream-Completed");
            System.out.println("流传输完成,Trailer值:" + completed);
        } finally {
            client.close();
        }
    }
}

内容的提问来源于stack exchange,提问作者Josef Mader-Kreiner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 09:20:21