流式 RPC
流式 RPC
了解如何使用流式 RPC。
在某些场景下,客户端或服务端需要发送海量数据,这些数据可能会随时间不断增长,或者体积过大而无法放入 RPC 附件中。例如,分布式系统中不同节点之间传输的副本或快照。虽然我们可以把数据分段,通过客户端与服务端之间的多次 RPC 来发送,但这会带来以下问题:
- 如果这些 RPC 是并行发送的,则无法保证接收端的数据顺序,从而导致重组代码变得复杂。
- 如果这些 RPC 是串行发送的,我们不得不忍受每次 RPC 的网络往返时延(RTT)与处理时间的叠加,这尤其难以预测。
为了使大数据包能够像流一样在客户端与服务端之间传输,我们提供了一种新的通信模型:Streaming RPC。Streaming RPC 允许用户建立 Stream,即客户端与服务之间的一条用户态连接。多条 Stream 可以同时共享同一条 TCP 连接。Stream 上的基本传输单元是 message。因此,发送方可以持续地向 Stream 写入消息,而接收方则可以按发送顺序将它们读出。
流式 RPC 确保/提供:
- 接收方的消息顺序与发送方完全一致
- 消息的边界
- 全双工
- 流量控制
- 超时通知
我们暂不支持自动切分大消息,因此在单个 TCP 连接上运行多个 Stream 可能导致队头阻塞(Head-of-line blocking)问题。在我们提供自动分段功能之前,请避免在单条消息中放入大量数据。
有关示例,请参阅 example/streaming_echo_c++。
创建流
目前流仅由客户端建立。客户端中会创建一个新的 Stream 对象,并用它通过 baidu_std 协议向指定服务发起 RPC。服务可以通过无错误地响应请求来接受该流,因此当客户端成功收到响应时,Stream 即告创建成功。此过程中的任何错误都会导致 RPC 失败,进而导致 Stream 创建失败。以 Linux 环境为例,客户端先创建一个 socket(创建 Stream),然后通过 connect 尝试与远端建立连接(通过 RPC 建立 Stream)。最终,当远端 accept 该请求后,流即创建完成。
如果客户端尝试与不支持流式 RPC 的服务器建立流,将会始终返回失败。
在代码中,我们使用 StreamId 来表示 Stream,它是读取、写入和关闭 Stream 时需要传入的关键 ID。
struct StreamOptions
// The max size of unconsumed data allowed at remote side.
// If |max_buf_size| <= 0, there's no limit of buf size
// default: 2097152 (2M)
int max_buf_size;
// Notify user when there's no data for at least |idle_timeout_ms|
// milliseconds since the last time that on_received_messages or on_idle_timeout
// finished.
// default: -1
long idle_timeout_ms;
// Maximum messages in batch passed to handler->on_received_messages
// default: 128
size_t messages_in_batch;
// Handle input message, if handler is NULL, the remote side is not allowd to
// write any message, who will get EBADF on writting
// default: NULL
StreamInputHandler* handler;
};
// [Called at the client side]
// Create a Stream at client-side along with the |cntl|, which will be connected
// when receiving the response with a Stream from server-side. If |options| is
// NULL, the Stream will be created with default options
// Return 0 on success, -1 otherwise
int StreamCreate(StreamId* request_stream, Controller &cntl, const StreamOptions* options);接受流
如果在 RPC 的请求中附带了一个 Stream,服务端可以通过 StreamAccept 接受该 Stream。成功时,此函数会将创建好的 Stream 填入 response_stream,可用于向客户端发送消息。
// [Called at the server side]
// Accept the Stream. If client didn't create a Stream with the request
// (cntl.has_remote_stream() returns false), this method would fail.
// Return 0 on success, -1 otherwise.
int StreamAccept(StreamId* response_stream, Controller &cntl, const StreamOptions* options);从流中读取
创建/接受一个 Stream 后,你可以用自己实现的 StreamInputHandler 填充 StreamOptions 中的 hander。之后,当流收到数据、被对端关闭或达到空闲超时时间时,你会收到通知。
class StreamInputHandler {
public:
// Callback when stream receives data
virtual int on_received_messages(StreamId id, butil::IOBuf *const messages[], size_t size) = 0;
// Callback when there is no data for a long time on the stream
virtual void on_idle_timeout(StreamId id) = 0;
// Callback when stream is closed by the other end
virtual void on_closed(StreamId id) = 0;
};对
on_received_message的首次调用在客户端,如果创建过程是同步的,
on_received_message会在阻塞式 RPC 返回时被调用。如果是异步的,on_received_message要等到done->Run()结束后才会被调用。在服务端,
on_received_message会在done->Run()结束后被调用。
向流中写入
// Write |message| into |stream_id|. The remote-side handler will received the
// message by the written order
// Returns 0 on success, errno otherwise
// Errno:
// - EAGAIN: |stream_id| is created with positive |max_buf_size| and buf size
// which the remote side hasn't consumed yet excceeds the number.
// - EINVAL: |stream_id| is invalied or has been closed
int StreamWrite(StreamId stream_id, const butil::IOBuf &message);流控
当未确认数据量达到上限时,发送端的 Write 操作会立即以 EAGAIN 失败。此时,你应当同步或异步地等待接收端消费数据。
// Wait util the pending buffer size is less than |max_buf_size| or error occurs
// Returns 0 on success, errno otherwise
// Errno:
// - ETIMEDOUT: when |due_time| is not NULL and time expired this
// - EINVAL: the Stream was close during waiting
int StreamWait(StreamId stream_id, const timespec* due_time);
// Async wait
void StreamWait(StreamId stream_id, const timespec *due_time,
void (*on_writable)(StreamId stream_id, void* arg, int error_code),
void *arg);关闭流
// Close |stream_id|, after this function is called:
// - All the following |StreamWrite| would fail
// - |StreamWait| wakes up immediately.
// - Both sides |on_closed| would be notifed after all the pending buffers have
// been received
// This function could be called multiple times without side-effects
int StreamClose(StreamId stream_id);最后修改于 2022 年 1 月 9 日:[基于 hugo 的 brpc 网站新版本 (94b25d711)]](https://github.com/apache/brpc-website/commit/94b25d7110944f3d4b2071b5188c691b23ffe3a9)
评论
登录后参与评论
KnowForge