RPC 换成 HTTP 后,发生了什么?自定义协议与 HTTP,怎么选?
前言
前面的 RPC 使用自定义二进制协议:
┌───────────────┬──────────────────┐
│ 自定义 Header │ 序列化后的 Body │
└───────────────┴──────────────────┘
Header 中包含:
魔数
版本
消息类型
requestId
Body 长度
序列化方式
客户端通过 Netty 发送数据,服务端使用自定义解码器完成:
拆帧
↓
解析 Header
↓
反序列化 Body
↓
执行远程方法
现在我们准备做一件看似简单的事情:
不再使用自定义 Header
改用 HTTP 传输 RpcRequest
调用者仍然这样写:
CarService carService =
rpcProxyFactory.create(CarService.class);
Person person =
carService.findPerson("Tom", 20);
但是底层链路变成:
动态代理
↓
RpcInvocation
↓
HTTP POST
↓
服务端处理
↓
HTTP Response
↓
返回 Person
于是新的问题出现了:
换成 HTTP 后,动态代理需要重写吗?
RpcInvocation 应该放在 URL、Header 还是 Body 中?
HTTP 是无状态协议,是不是每个请求都必须新建 TCP 连接?
自定义协议有 requestId,HTTP 为什么没有?
多个请求能不能复用同一条 HTTP/1.1 连接?
HTTP/1.1 上的响应能不能乱序返回?
没有设置
Content-Length,为什么服务端读取不到请求体?
这到底是 TCP 粘包,还是 HTTP 消息格式错误?
HttpServerCodec与自定义 RPC 解码器有什么区别?
HttpObjectAggregator为什么能直接得到FullHttpRequest?
使用
HttpURLConnection和使用 Netty,协议层有什么不同?
Provider 换成 Jetty 或 Tomcat 后,Consumer 是否还要修改?
这些问题最终指向一个核心:
RPC 的调用模型与传输协议应该如何解耦?
一、哪些不该变?
将自定义协议换成 HTTP,并不代表整个 RPC 都要推翻重写。
原有调用链中,可以保持不变的部分包括:
服务接口
动态代理
RpcInvocation
本地与远程路由
服务发现
负载均衡
服务注册表
方法反射调用
返回值与异常模型
真正需要替换的是:
传输层
+
协议编解码层
原来的结构:
RpcInvocation
↓
RpcRequest
↓
自定义 Header + Body
↓
Netty TCP
替换后:
RpcInvocation
↓
RpcRequest
↓
HTTP Request
↓
HTTP Client
服务端则从:
Netty TCP
↓
自定义 FrameDecoder
↓
RpcRequest
变成:
HTTP Server
↓
HTTP 协议解码
↓
HTTP Body
↓
RpcRequest
因此,分层正确的 RPC 只需要替换:
RpcTransport
例如:
public interface RpcTransport {
CompletionStage<RpcResponse> request(
Endpoint endpoint,
RpcRequest request
);
}
可以有两个实现:
NettyBinaryTransport
HttpRpcTransport
代理层并不需要知道使用的是哪一种。
二、HTTP 无状态吗?
HTTP 确实被定义为无状态的应用层请求—响应协议。
但“无状态”的准确含义是:
每个请求消息的语义都可以被独立理解,服务器不应该仅根据两个请求使用了同一条连接,就假设它们属于同一个用户或同一次业务会话。
它描述的是请求语义,而不是 TCP 连接必须用完即关。
因此:
HTTP 无状态
不能推出:
每个请求必须建立一条新 TCP 连接
也不能推出:
HTTP 不能保存登录状态
登录会话可以通过:
Cookie
Token
Session ID
Authorization
等应用层信息表达。
RPC 调用关联也可以通过:
requestId
traceId
业务请求号
表达。
这些字段增加的是:
应用关联状态
而不是改变 HTTP 本身的无状态语义。
三、无状态不等于短连接
HTTP/1.1 默认使用持久连接,一条连接可以连续承载多个请求和响应。
只要没有发送:
Connection: close
并且消息都有明确的结束边界,连接通常可以在当前响应后继续使用。
所以:
一个 HttpURLConnection 对象
对应一次请求—响应交换
并不等于:
每次一定创建一条全新的 TCP 连接
同样:
创建一个新的 FullHttpRequest
也不等于:
必须创建一个新的 Channel
真正决定连接模型的是客户端实现:
每次请求新建 Channel
连接池复用 Channel
单连接串行请求
多连接并行请求
HTTP/2 单连接多路复用
四、三种状态别混
讨论“有状态”时,至少要区分三个层次。
1. 会话状态
例如:
用户是否登录
购物车内容
事务上下文
订阅关系
2. 连接状态
例如:
TCP 序号
TLS 会话
HTTP/2 Stream
流量控制窗口
3. 调用状态
例如:
requestId → CompletableFuture
在 RPC 请求中增加 requestId,表示客户端维护了未完成调用的关联状态:
1001 → Future A
1002 → Future B
1003 → Future C
这并不意味着:
HTTP 已经从无状态协议
变成了有状态协议
更准确的说法是:
RPC 在 HTTP 消息之上增加了自己的请求关联机制。
五、货物与火车
资料中将业务数据称为“货物”,将底层协议称为“小火车”。
这个类比非常适合解释分层。
1. 货物
RPC 真正想传输的是:
public record RpcInvocation(
String serviceName,
String methodName,
String[] parameterTypeNames,
Object[] arguments
) {
}
它描述:
调用哪个服务
调用哪个方法
参数类型是什么
参数值是什么
它不应该知道:
HTTP Header 怎么写
自定义协议魔数是多少
Netty Channel 在哪里
Content-Length 怎样设置
2. 火车
传输载体可以是:
自定义二进制协议
HTTP/1.1
HTTP/2
其他消息协议
自定义协议:
┌───────────────┬──────────────────┐
│ RpcFrameHeader│ RpcRequest Body │
└───────────────┴──────────────────┘
HTTP:
POST /rpc HTTP/1.1
Content-Type: application/octet-stream
Content-Length: 324
<序列化后的 RpcRequest>
3. 分层结果
RpcInvocation
↓
RpcRequest
↓
RpcSerializer
↓
byte[]
↓
┌─────────────────┬──────────────────┐
│ 自定义协议 Transport │
│ 或 HTTP Transport │
└─────────────────┴──────────────────┘
更换火车,不应该改变货物。
六、Transport 怎么拆?
可以定义统一接口:
public interface RpcTransport {
CompletionStage<RpcResponse> request(
Endpoint endpoint,
RpcRequest request
);
}
远程调用器只依赖它:
public final class RemoteInvoker
implements RpcInvoker {
private final ServiceDiscovery discovery;
private final LoadBalancer loadBalancer;
private final RpcTransport transport;
private final AtomicLong requestIds =
new AtomicLong();
@Override
public CompletionStage<Object> invoke(
RpcInvocation invocation
) {
List<Endpoint> endpoints =
discovery.discover(
invocation.serviceKey()
);
Endpoint endpoint =
loadBalancer.select(
endpoints,
invocation
);
RpcRequest request =
new RpcRequest(
requestIds.incrementAndGet(),
invocation
);
return transport.request(endpoint, request)
.thenApply(response -> {
if (!response.success()) {
throw new RpcRemoteException(
response.errorType(),
response.errorMessage()
);
}
return response.result();
});
}
}
切换实现:
RpcTransport transport =
new NettyBinaryTransport(...);
或者:
RpcTransport transport =
new NettyHttpTransport(...);
甚至:
RpcTransport transport =
new UrlConnectionTransport(...);
代理层完全不需要变化。
七、先用阻塞客户端
为了理解 HTTP 链路,可以先使用 HttpURLConnection。
它的优势不是性能,而是结构直接:
创建请求
↓
写入 Body
↓
等待响应
↓
读取响应 Body
示例:
public final class UrlConnectionTransport
implements RpcTransport {
private final RpcSerializer serializer;
private final Duration connectTimeout;
private final Duration readTimeout;
public UrlConnectionTransport(
RpcSerializer serializer,
Duration connectTimeout,
Duration readTimeout
) {
this.serializer = serializer;
this.connectTimeout = connectTimeout;
this.readTimeout = readTimeout;
}
@Override
public CompletionStage<RpcResponse> request(
Endpoint endpoint,
RpcRequest request
) {
try {
RpcResponse response =
doRequest(endpoint, request);
return CompletableFuture.completedFuture(
response
);
} catch (Throwable throwable) {
return CompletableFuture.failedFuture(
throwable
);
}
}
private RpcResponse doRequest(
Endpoint endpoint,
RpcRequest request
) throws Exception {
byte[] requestBody =
serializer.serialize(request);
URI uri = URI.create(
"http://"
+ endpoint.host()
+ ':'
+ endpoint.port()
+ "/rpc"
);
HttpURLConnection connection =
(HttpURLConnection)
uri.toURL().openConnection();
connection.setRequestMethod("POST");
connection.setDoOutput(true);
connection.setDoInput(true);
connection.setConnectTimeout(
Math.toIntExact(
connectTimeout.toMillis()
)
);
connection.setReadTimeout(
Math.toIntExact(
readTimeout.toMillis()
)
);
connection.setRequestProperty(
"Content-Type",
"application/octet-stream"
);
connection.setFixedLengthStreamingMode(
requestBody.length
);
try {
try (OutputStream output =
connection.getOutputStream()) {
output.write(requestBody);
}
int status =
connection.getResponseCode();
if (status < 200 || status >= 300) {
throw new RpcTransportException(
"HTTP status: " + status
);
}
try (InputStream input =
connection.getInputStream()) {
byte[] responseBody =
input.readAllBytes();
return serializer.deserialize(
responseBody,
RpcResponse.class
);
}
} finally {
connection.disconnect();
}
}
}
URLConnection 在真正建立连接前可以先配置输入、输出和超时选项;依赖连接的操作会在必要时隐式建立连接。HttpURLConnection.getResponseCode() 则读取 HTTP 响应状态码。
阻塞在哪里?
这段代码是同步阻塞的:
业务线程
↓
写入 HTTP 请求
↓
等待响应状态码
↓
等待响应 Body
↓
返回 RpcResponse
即使接口返回:
CompletionStage<RpcResponse>
如果内部先同步执行完整请求,再创建一个已完成的 Future,它仍然没有释放调用线程。
真正的异步传输,需要让网络事件由 EventLoop 推进。
八、服务端怎么解析?
Netty HTTP 服务端常见 Pipeline:
pipeline.addLast(
new HttpServerCodec(),
new HttpObjectAggregator(MAX_CONTENT_LENGTH),
new HttpRpcServerHandler(...)
);
1. HttpServerCodec
HttpServerCodec 组合了:
HttpRequestDecoder
+
HttpResponseEncoder
也就是:
入站:
HTTP 字节 → HttpRequest / HttpContent
出站:
HttpResponse → HTTP 字节
它承担的是 HTTP 协议编解码,而不理解 RPC 业务对象。
2. HttpObjectAggregator
HTTP 请求体可能被拆成:
HttpRequest
HttpContent
HttpContent
LastHttpContent
HttpObjectAggregator 会把 HttpMessage 及其后续 HttpContent 聚合成一个:
FullHttpRequest
响应方向则可以聚合成:
FullHttpResponse
这样业务 Handler 不必自己保存多个 HTTP 内容片段。
Pipeline 变成:
TCP 字节
↓
HttpServerCodec
↓
HttpRequest + HttpContent
↓
HttpObjectAggregator
↓
FullHttpRequest
↓
HttpRpcServerHandler
九、服务端完整示例
public final class HttpRpcServer {
private static final int MAX_CONTENT_LENGTH =
8 * 1024 * 1024;
private final EventLoopGroup bossGroup =
new NioEventLoopGroup(1);
private final EventLoopGroup workerGroup =
new NioEventLoopGroup();
private final EventExecutorGroup businessGroup =
new DefaultEventExecutorGroup(16);
private final RpcSerializer serializer;
private final RequestDispatcher dispatcher;
public HttpRpcServer(
RpcSerializer serializer,
RequestDispatcher dispatcher
) {
this.serializer = serializer;
this.dispatcher = dispatcher;
}
public ChannelFuture start(int port) {
ServerBootstrap bootstrap =
new ServerBootstrap();
return bootstrap
.group(bossGroup, workerGroup)
.channel(
NioServerSocketChannel.class
)
.childHandler(
new ChannelInitializer<
SocketChannel>() {
@Override
protected void initChannel(
SocketChannel channel
) {
ChannelPipeline pipeline =
channel.pipeline();
pipeline.addLast(
new HttpServerCodec()
);
pipeline.addLast(
new HttpObjectAggregator(
MAX_CONTENT_LENGTH
)
);
pipeline.addLast(
businessGroup,
new HttpRpcServerHandler(
serializer,
dispatcher
)
);
}
}
)
.bind(port);
}
}
Handler:
public final class HttpRpcServerHandler
extends SimpleChannelInboundHandler<
FullHttpRequest> {
private static final String RPC_PATH =
"/rpc";
private final RpcSerializer serializer;
private final RequestDispatcher dispatcher;
public HttpRpcServerHandler(
RpcSerializer serializer,
RequestDispatcher dispatcher
) {
this.serializer = serializer;
this.dispatcher = dispatcher;
}
@Override
protected void channelRead0(
ChannelHandlerContext context,
FullHttpRequest request
) {
if (request.method() != HttpMethod.POST
|| !RPC_PATH.equals(request.uri())) {
writeStatus(
context,
request,
HttpResponseStatus.NOT_FOUND
);
return;
}
byte[] requestBytes =
new byte[
request.content()
.readableBytes()
];
request.content().readBytes(
requestBytes
);
boolean keepAlive =
HttpUtil.isKeepAlive(request);
RpcRequest rpcRequest;
try {
rpcRequest =
serializer.deserialize(
requestBytes,
RpcRequest.class
);
} catch (Throwable throwable) {
writeStatus(
context,
request,
HttpResponseStatus.BAD_REQUEST
);
return;
}
dispatcher.dispatch(rpcRequest)
.whenComplete(
(rpcResponse, throwable) -> {
RpcResponse response =
throwable == null
? rpcResponse
: RpcResponse.failure(
rpcRequest
.requestId(),
throwable
);
writeResponse(
context,
keepAlive,
response
);
}
);
}
private void writeResponse(
ChannelHandlerContext context,
boolean keepAlive,
RpcResponse rpcResponse
) {
try {
byte[] responseBytes =
serializer.serialize(
rpcResponse
);
ByteBuf content =
Unpooled.wrappedBuffer(
responseBytes
);
FullHttpResponse response =
new DefaultFullHttpResponse(
HttpVersion.HTTP_1_1,
HttpResponseStatus.OK,
content
);
response.headers().set(
HttpHeaderNames.CONTENT_TYPE,
"application/octet-stream"
);
HttpUtil.setContentLength(
response,
content.readableBytes()
);
if (keepAlive) {
response.headers().set(
HttpHeaderNames.CONNECTION,
HttpHeaderValues.KEEP_ALIVE
);
context.writeAndFlush(response);
return;
}
response.headers().set(
HttpHeaderNames.CONNECTION,
HttpHeaderValues.CLOSE
);
context.writeAndFlush(response)
.addListener(
ChannelFutureListener.CLOSE
);
} catch (Throwable throwable) {
context.close();
}
}
private void writeStatus(
ChannelHandlerContext context,
FullHttpRequest request,
HttpResponseStatus status
) {
boolean keepAlive =
HttpUtil.isKeepAlive(request);
FullHttpResponse response =
new DefaultFullHttpResponse(
HttpVersion.HTTP_1_1,
status
);
HttpUtil.setContentLength(
response,
0
);
ChannelFuture future =
context.writeAndFlush(response);
if (!keepAlive) {
future.addListener(
ChannelFutureListener.CLOSE
);
}
}
@Override
public void exceptionCaught(
ChannelHandlerContext context,
Throwable cause
) {
cause.printStackTrace();
context.close();
}
}
十、Aggregator 解决什么?
HttpObjectAggregator 经常被描述为解决“HTTP 粘包”。
这个说法不够准确。
它真正解决的是:
把已经被 HTTP 解码器识别的
HttpMessage 与 HttpContent
聚合成完整 HTTP 消息
也就是说:
TCP 字节边界
↓
HttpServerCodec 解析 HTTP 消息语法
↓
HttpObjectAggregator 聚合内容片段
它不是直接在原始 TCP 字节中猜测:
哪一段属于请求 A
哪一段属于请求 B
HTTP 消息边界首先由 HTTP 协议字段决定,包括:
Content-Length
Transfer-Encoding
请求方法
响应状态
连接关闭语义
Aggregator 只能在消息格式本身正确的前提下完成聚合。
十一、Content-Length 是什么?
假设发送一个 POST 请求:
POST /rpc HTTP/1.1
Host: localhost:9090
Content-Type: application/octet-stream
<324 字节请求体>
如果既没有:
Content-Length: 324
也没有:
Transfer-Encoding: chunked
HTTP/1.1 接收方无法把这 324 字节识别为该请求的 Body。
对于请求消息,如果没有其他规则定义 Body 长度,并且同时缺少有效的 Content-Length 与 Transfer-Encoding,协议会把请求体长度视为 0;服务器也可以对缺少长度信息的请求返回 411 Length Required。
所以问题不是:
TCP 粘包导致 HTTP 读不到 Body
而是:
HTTP 消息没有正确声明 Body 边界
正确设置
客户端:
HttpUtil.setContentLength(
request,
content.readableBytes()
);
或者:
request.headers().set(
HttpHeaderNames.CONTENT_LENGTH,
content.readableBytes()
);
服务端响应同样需要设置:
HttpUtil.setContentLength(
response,
content.readableBytes()
);
Netty 的 HttpUtil.setContentLength() 就是用于设置 HTTP 消息的 Content-Length Header。
十二、Chunked 可以吗?
HTTP/1.1 并不要求所有 Body 都必须通过 Content-Length 定界。
也可以使用:
Transfer-Encoding: chunked
例如:
4\r\n
Wiki\r\n
5\r\n
pedia\r\n
0\r\n
\r\n
接收方根据每个 Chunk 的长度完成消息解析。
因此,准确的规则是:
已知请求体长度
优先使用 Content-Length
流式或长度未知
可以使用 chunked
HTTP/1.1 发送带 Body 的请求时,必须提供有效的 Content-Length,或者使用 Chunked Transfer Coding。
使用 HttpObjectAggregator 时,无论内容最初来自固定长度 Body 还是多个 Chunk,它都可以在上限内聚合成 FullHttpRequest。
十三、客户端怎么写?
Netty HTTP 客户端 Pipeline:
HttpClientCodec
↓
HttpObjectAggregator
↓
HttpRpcClientHandler
HttpClientCodec 组合了:
HttpRequestEncoder
+
HttpResponseDecoder
用于完成客户端方向的 HTTP 请求编码和响应解码。
客户端骨架:
public final class NettyHttpTransport
implements RpcTransport, AutoCloseable {
private static final int MAX_CONTENT_LENGTH =
8 * 1024 * 1024;
private final EventLoopGroup group =
new NioEventLoopGroup();
private final RpcSerializer serializer;
private final ConcurrentMap<
Endpoint,
CompletableFuture<Channel>
> channels = new ConcurrentHashMap<>();
private final ConcurrentMap<
Long,
CompletableFuture<RpcResponse>
> pendingRequests =
new ConcurrentHashMap<>();
private final Bootstrap bootstrap;
public NettyHttpTransport(
RpcSerializer serializer
) {
this.serializer = serializer;
this.bootstrap = new Bootstrap()
.group(group)
.channel(NioSocketChannel.class)
.handler(
new ChannelInitializer<
SocketChannel>() {
@Override
protected void initChannel(
SocketChannel channel
) {
channel.pipeline().addLast(
new HttpClientCodec(),
new HttpObjectAggregator(
MAX_CONTENT_LENGTH
),
new HttpRpcClientHandler(
serializer,
pendingRequests
)
);
}
}
);
}
@Override
public CompletionStage<RpcResponse> request(
Endpoint endpoint,
RpcRequest request
) {
return getChannel(endpoint)
.thenCompose(channel ->
send(channel, endpoint, request)
);
}
private CompletionStage<Channel> getChannel(
Endpoint endpoint
) {
return channels.computeIfAbsent(
endpoint,
this::connect
);
}
private CompletableFuture<Channel> connect(
Endpoint endpoint
) {
CompletableFuture<Channel> result =
new CompletableFuture<>();
bootstrap.connect(
endpoint.host(),
endpoint.port()
)
.addListener(future -> {
if (!future.isSuccess()) {
channels.remove(
endpoint,
result
);
result.completeExceptionally(
future.cause()
);
return;
}
Channel channel =
((ChannelFuture) future)
.channel();
result.complete(channel);
});
return result;
}
@Override
public void close() {
group.shutdownGracefully();
}
}
Netty 的 Channel I/O 操作是异步的,调用后会立即返回 ChannelFuture,由 Future 表达连接或写入的最终结果。
十四、请求怎么构造?
private CompletionStage<RpcResponse> send(
Channel channel,
Endpoint endpoint,
RpcRequest rpcRequest
) {
CompletableFuture<RpcResponse> result =
new CompletableFuture<>();
pendingRequests.put(
rpcRequest.requestId(),
result
);
byte[] requestBytes;
try {
requestBytes =
serializer.serialize(
rpcRequest
);
} catch (Throwable throwable) {
pendingRequests.remove(
rpcRequest.requestId()
);
return CompletableFuture.failedFuture(
throwable
);
}
ByteBuf content =
Unpooled.wrappedBuffer(
requestBytes
);
FullHttpRequest request =
new DefaultFullHttpRequest(
HttpVersion.HTTP_1_1,
HttpMethod.POST,
"/rpc",
content
);
request.headers().set(
HttpHeaderNames.HOST,
endpoint.host()
+ ':'
+ endpoint.port()
);
request.headers().set(
HttpHeaderNames.CONTENT_TYPE,
"application/octet-stream"
);
request.headers().set(
HttpHeaderNames.CONNECTION,
HttpHeaderValues.KEEP_ALIVE
);
HttpUtil.setContentLength(
request,
content.readableBytes()
);
channel.writeAndFlush(request)
.addListener(writeFuture -> {
if (writeFuture.isSuccess()) {
return;
}
CompletableFuture<RpcResponse> removed =
pendingRequests.remove(
rpcRequest.requestId()
);
if (removed != null) {
removed.completeExceptionally(
writeFuture.cause()
);
}
});
return result;
}
这里仍然要区分两个 Future:
ChannelFuture
表示 HTTP 请求写操作是否成功
CompletableFuture<RpcResponse>
表示远程方法是否返回结果
HTTP 请求成功写出,不代表:
服务端已经完成反序列化
业务方法已经执行
响应已经返回
十五、响应怎么匹配?
客户端 Handler:
public final class HttpRpcClientHandler
extends SimpleChannelInboundHandler<
FullHttpResponse> {
private final RpcSerializer serializer;
private final ConcurrentMap<
Long,
CompletableFuture<RpcResponse>
> pendingRequests;
public HttpRpcClientHandler(
RpcSerializer serializer,
ConcurrentMap<
Long,
CompletableFuture<RpcResponse>
> pendingRequests
) {
this.serializer = serializer;
this.pendingRequests = pendingRequests;
}
@Override
protected void channelRead0(
ChannelHandlerContext context,
FullHttpResponse response
) {
if (response.status().code() < 200
|| response.status().code() >= 300) {
failOldestRequest(
new RpcTransportException(
"HTTP status: "
+ response.status()
)
);
return;
}
byte[] responseBytes =
new byte[
response.content()
.readableBytes()
];
response.content().readBytes(
responseBytes
);
RpcResponse rpcResponse;
try {
rpcResponse =
serializer.deserialize(
responseBytes,
RpcResponse.class
);
} catch (Throwable throwable) {
failOldestRequest(throwable);
return;
}
CompletableFuture<RpcResponse> future =
pendingRequests.remove(
rpcResponse.requestId()
);
if (future != null) {
future.complete(rpcResponse);
}
}
private void failOldestRequest(
Throwable throwable
) {
pendingRequests.entrySet()
.stream()
.findFirst()
.ifPresent(entry -> {
if (pendingRequests.remove(
entry.getKey(),
entry.getValue()
)) {
entry.getValue()
.completeExceptionally(
throwable
);
}
});
}
}
这段代码展示了 requestId 的基本作用,但在 HTTP/1.1 中还存在一个关键限制。
十六、HTTP/1.1 能乱序吗?
自定义 RPC 协议可以设计成:
Request 101
Request 102
Request 103
Response 103
Response 101
Response 102
只要每条响应携带 requestId,客户端就能完成匹配。
但 HTTP/1.1 不包含用于关联请求与响应的协议级请求标识。
如果同一条连接上存在多个未完成请求,客户端按照发送顺序维护请求队列,并把收到的响应依次关联给最早尚未收到最终响应的请求。
HTTP/1.1 虽然允许 Pipelining:
先连续发送 Request A、B、C
不等待每个响应
但服务端必须按照请求到达顺序发送对应响应:
Request A → Response A
Request B → Response B
Request C → Response C
即使 B 的业务先执行完成,也不能先于 A 返回。
所以:
在 HTTP/1.1 上增加业务 requestId,并不能自动允许响应在同一连接中任意乱序发送。
十七、并发怎么做?
使用 HTTP/1.1 作为 RPC 传输时,有三种常见选择。
1. 单连接串行
发送请求 A
↓
等待响应 A
↓
发送请求 B
优点:
实现简单
响应顺序明确
缺点:
慢请求会阻塞后续请求
2. 多连接并发
Channel 1 → 请求 A
Channel 2 → 请求 B
Channel 3 → 请求 C
每条连接同时只保留一个未完成请求。
优点:
容易实现
支持并行
缺点:
连接数量增加
资源消耗更高
3. HTTP/1.1 Pipelining
单连接连续发送多个请求
服务端按顺序返回响应
这仍然存在响应层的队头阻塞:
请求 A 很慢
↓
请求 B 已完成
↓
但响应 B 必须等待响应 A
HTTP/1.1 规范也指出,多连接经常被用来缓解这种队头阻塞,但每条连接都会消耗额外服务端资源。
对于教学版 RPC,更容易写对的模型是:
HTTP/1.1
+
连接池
+
每条连接一个未完成请求
十八、HTTP/2 有什么不同?
HTTP/2 在一条连接中引入多个独立 Stream:
HTTP/2 Connection
├── Stream 1
├── Stream 3
├── Stream 5
└── Stream 7
不同 Stream 的 Frame 可以在一条连接中交错传输,每个 Stream 由独立整数 ID 标识。
因此:
Request A → Stream 1
Request B → Stream 3
Request C → Stream 5
响应可以根据各自 Stream 独立推进:
Response B → Stream 3
Response C → Stream 5
Response A → Stream 1
这才更接近前面自定义 RPC 协议的:
单连接
+
多请求并发
+
乱序完成
模型。
HTTP/2 已经具有协议级 Stream ID。
业务 requestId 仍然可以用于:
分布式追踪
业务幂等
跨连接重试
日志关联
框架内部调用标识
但不再需要用它替代 HTTP/2 的底层 Stream 关联。
十九、HTTP Header 放什么?
RPC 数据可以全部放在 Body:
RpcRequest
├── requestId
├── serviceName
├── methodName
├── parameterTypes
└── arguments
HTTP Header 只保存传输元信息:
Content-Type: application/octet-stream
Content-Length: 324
Connection: keep-alive
也可以增加:
Rpc-Version: 1
Rpc-Serializer: java
Rpc-Request-Id: 1001
Traceparent: ...
但不要把完整业务参数放入自定义 Header。
更合理的分工是:
HTTP Header
描述 HTTP 消息和跨层元数据
HTTP Body
保存 RpcRequest 或 RpcResponse
requestId 可以同时出现在:
HTTP Header
+
RpcRequest Body
但必须定义哪一个是权威来源,并检查二者是否一致,否则会产生歧义。
简单实现中,只放在 Body 即可。
二十、Provider 可以替换吗?
使用 HTTP 的一个重要优势是:
客户端与服务端只要遵守 HTTP
就不必使用相同的 I/O 框架
Consumer 可以是:
HttpURLConnection
Java HttpClient
Netty HTTP Client
其他 HTTP 客户端
Provider 可以是:
Netty HTTP Server
Jetty
Tomcat
Servlet 容器
其他 Web Server
组合关系:
| Consumer | Provider |
|---|---|
HttpURLConnection | Netty |
HttpURLConnection | Jetty |
| Netty HTTP Client | Netty |
| Netty HTTP Client | Jetty |
| Java HTTP Client | Tomcat |
只要双方遵守相同约定:
POST /rpc
Content-Type
Body 序列化格式
RpcRequest 字段
RpcResponse 字段
错误码约定
调用层就可以互通。
这正是标准协议带来的生态价值。
二十一、Servlet 怎么接?
基于 Servlet 的 Provider 可以这样处理:
public final class RpcServlet
extends HttpServlet {
private final RpcSerializer serializer;
private final RequestDispatcher dispatcher;
public RpcServlet(
RpcSerializer serializer,
RequestDispatcher dispatcher
) {
this.serializer = serializer;
this.dispatcher = dispatcher;
}
@Override
protected void doPost(
HttpServletRequest request,
HttpServletResponse response
) throws IOException {
byte[] requestBytes;
try (InputStream input =
request.getInputStream()) {
requestBytes =
input.readAllBytes();
}
RpcRequest rpcRequest;
try {
rpcRequest =
serializer.deserialize(
requestBytes,
RpcRequest.class
);
} catch (Throwable throwable) {
response.sendError(
HttpServletResponse
.SC_BAD_REQUEST
);
return;
}
RpcResponse rpcResponse;
try {
rpcResponse =
dispatcher.dispatch(rpcRequest)
.toCompletableFuture()
.join();
} catch (Throwable throwable) {
rpcResponse =
RpcResponse.failure(
rpcRequest.requestId(),
throwable
);
}
byte[] responseBytes;
try {
responseBytes =
serializer.serialize(
rpcResponse
);
} catch (Throwable throwable) {
response.sendError(
HttpServletResponse
.SC_INTERNAL_SERVER_ERROR
);
return;
}
response.setStatus(
HttpServletResponse.SC_OK
);
response.setContentType(
"application/octet-stream"
);
response.setContentLength(
responseBytes.length
);
response.getOutputStream()
.write(responseBytes);
}
}
不过这里的:
.join();
会阻塞当前 Servlet 请求线程。
要实现完整异步链路,需要结合 Servlet 异步处理或响应式 Web 容器,而不是仅让底层 Transport 返回 Future。
二十二、局部异步有用吗?
即使最外层接口是同步的:
Person findPerson(
String name,
int age
);
底层使用 Netty 异步传输仍然有价值。
因为:
业务调用线程
可以等待 Future
Netty EventLoop
不需要阻塞等待网络数据
这两个线程不是同一类资源。
因此:
最外层同步
不等于:
底层异步毫无意义
不过,如果业务链路每一层最终都立刻调用:
future.get();
就会保留大量等待线程。
真正的全链路异步通常要求接口本身返回:
CompletionStage<Person>
并让上层容器也支持异步响应。
二十三、协议怎么选择?
自定义协议
适合:
内部服务通信
追求更低协议开销
需要单连接乱序多路复用
需要精细控制协议字段
通信双方完全可控
优势:
Header 紧凑
编解码路径短
可自由设计 requestId
容易针对 RPC 优化
代价:
需要自己实现协议
需要处理版本兼容
需要维护客户端与服务端
通用工具支持较少
HTTP/1.1
适合:
跨语言
跨容器
对接现有 Web 基础设施
需要代理、网关和可观测工具
性能要求适中
优势:
生态成熟
调试方便
容器选择多
网络设备兼容性好
限制:
Header 开销更大
同连接响应要求有序
并发通常依赖连接池
HTTP/2
适合:
单连接并发请求
需要 Stream 多路复用
跨语言 RPC
希望复用 HTTP 生态
优势:
多路复用
协议级 Stream ID
头部压缩
单连接并发
代价:
实现和排查更复杂
需要 HTTP/2 客户端与服务端支持
仍需设计 RPC Body 和错误模型
二十四、常见误区
1. HTTP 无状态,所以不能实现 RPC
错误。
RPC 调用信息可以放在 HTTP 请求中,响应可以通过 HTTP 返回。
2. HTTP 无状态,所以每个请求必须新建连接
错误。
HTTP/1.1 默认支持持久连接,一条连接可以承载多个请求与响应。
3. 加上 requestId,HTTP 就变成了有状态协议
不准确。
增加的是 RPC 调用关联状态,不是改变 HTTP 的语义模型。
4. requestId 能让 HTTP/1.1 响应乱序
错误。
同一 HTTP/1.1 连接上的响应必须按照请求顺序关联和发送。
5. HTTP/1.1 不能复用连接
错误。
它默认使用持久连接。
6. HTTP/1.1 复用连接就等于 HTTP/2
错误。
HTTP/1.1 的请求与响应仍受顺序约束;HTTP/2 使用多个独立 Stream 实现真正的单连接并发多路复用。
7. 没有 Content-Length 就是 TCP 粘包
错误。
这是 HTTP 消息定界信息缺失,不能简单归因于 TCP 粘包。
8. HTTP Body 必须使用 Content-Length
不完整。
也可以使用 HTTP/1.1 Chunked Transfer Coding。
9. HttpServerCodec 会生成 FullHttpRequest
错误。
它负责 HTTP 请求解码与响应编码;HttpObjectAggregator 才负责把多个 HTTP 对象聚合成 FullHttpRequest。
10. Aggregator 直接解决 TCP 粘包
不准确。
它聚合的是已经经过 HTTP 解码的 HttpMessage 与 HttpContent。
11. 使用 FullHttpRequest 就不用限制大小
错误。
HttpObjectAggregator 必须配置最大内容长度,否则大请求会造成过高内存压力;超过上限时它可以返回 413 Request Entity Too Large。
12. 每个请求都应该创建 EventLoopGroup
错误。
EventLoopGroup 应被多个连接和请求共享。
13. 每个 HTTP 请求都应该创建 Bootstrap
错误。
Bootstrap 和 EventLoopGroup 可以长期复用,真正按目标维护的是 Channel 或连接池。
14. writeAndFlush 成功代表 RPC 成功
错误。
它只代表 Netty 写操作完成,不代表远程方法已经执行。
15. HTTP 状态码 200 就代表业务成功
不一定。
200 表示 HTTP 交换成功,RPC Body 中仍可能包含远程业务异常。
16. HTTP 状态码完全没用
错误。
协议解析失败、路径不存在、认证失败等传输或网关层错误,应该使用合适的 HTTP 状态码表达。
17. Provider 使用 Jetty,Consumer 就必须使用 Jetty
错误。
双方只需遵守相同 HTTP 与 RPC Body 约定。
18. 使用 HTTP 后不需要序列化
错误。
HTTP 解决消息传输格式,Java 对象仍需转换成 JSON、二进制或其他可传输形式。
19. HTTP 自动解决 RPC 幂等
错误。
HTTP 方法语义和业务方法幂等性不是同一件事,扣款等业务仍需设计业务请求号与去重机制。
20. 换成 HTTP 以后代理层也要重写
如果分层合理,代理层不应该感知传输协议变化。
总结
原来的自定义 RPC 链路是:
动态代理
↓
RpcInvocation
↓
RpcRequest
↓
自定义 Header + Body
↓
Netty TCP
换成 HTTP 后:
动态代理
↓
RpcInvocation
↓
RpcRequest
↓
HTTP POST Body
↓
HTTP Client
保持不变的是:
接口
代理
调用对象
服务发现
负载均衡
本地服务注册
方法执行
结果与异常模型
被替换的是:
协议层
传输层
客户端连接模型
服务端协议处理链
Netty HTTP 服务端链路是:
TCP 字节
↓
HttpServerCodec
↓
HttpRequest + HttpContent
↓
HttpObjectAggregator
↓
FullHttpRequest
↓
反序列化 RpcRequest
↓
RequestDispatcher
↓
RpcResponse
↓
FullHttpResponse
客户端链路是:
RpcRequest
↓
序列化
↓
DefaultFullHttpRequest
↓
HttpClientCodec
↓
网络发送
↓
FullHttpResponse
↓
RpcResponse
↓
完成 Future
HTTP 的无状态语义并不要求一请求一连接。
HTTP/1.1 默认支持持久连接,但同一连接上的响应需要按照请求顺序关联,因此不能仅凭业务 requestId 就实现任意乱序响应。
HTTP/2 则通过 Stream ID 在单条连接中提供并发多路复用,更接近自定义 RPC 的多请求并发模型。
Content-Length 也不是用来“优化粘包”的字段。
它负责告诉 HTTP 接收方:
当前消息的 Body 到哪里结束
如果 Body 长度未知,则可以使用 Chunked Transfer Coding。
最终,自定义协议与 HTTP 并不是谁绝对更好:
自定义协议
更紧凑、更可控、更适合内部高性能 RPC
HTTP
生态更成熟、更容易跨语言和跨容器
HTTP/2
在标准生态与多路复用之间取得平衡
最后,可以用一句话概括这次改造:
RPC 更换 HTTP 传输后,变化的应该只是协议与网络实现;代理、调用对象、服务路由和业务执行仍应保持稳定,这才是真正的分层解耦。
评论区