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

