lizining1231 opened a new issue, #3741: URL: https://github.com/apache/dubbo-go/issues/3741
## 摘要 Triple 流式客户端(`CallBidiStream` / `CallClientStream`)的典型用法是在**一条流上连续** **`Send`** **多条消息**,而当前实现里每次 `Send` 都把消息直接写进 `duplexHTTPCall` 底层的 `io.Pipe`,协议层与管道之间**没有任何缓冲层**。 `io.Pipe` 是无缓冲的跨 goroutine 同步交接:写端把切片交给读端后,要等读端回传已取走的字节数才能返回([pipe.go L76-L96](https://github.com/golang/go/blob/master/src/io/pipe.go#L76-L96))。这条链路上的读端不是业务代码,而是 `net/http` 从 `http.Request.Body` 取数据的 `writeRequestBody`;x/net/http2 又对每个 DATA 帧立即 `Flush`(源码注释原文 `TODO(bradfitz): this flush is for latency, not bandwidth`),于是一次下沉 `Write` 对应一次真实的 `write(2)`。 dubbo-go 默认编码是 gRPC wire,协议层把一条消息拆成 5 字节 prefix 与 body **两次** `Write`;Triple 编码的流式入口 `tripleUnaryMarshaler.write` 只有一次 body 写。两种编码每条消息的固定成本: | 环节 | gRPC wire(triple协议默认编码) | Triple 编码 | | -------------------------- | -------------------------- | --------------- | | 序列化 | 1 次 marshal | 1 次 marshal | | 协议层下沉 `Write` | **2 次**(5 B prefix + body) | **1 次**(仅 body) | | `io.Pipe` 跨 goroutine 同步交接 | **2 次** | **1 次** | | 传输层 `write(2)` | **2 次**(每 DATA 帧立即 Flush) | **1 次** | 后三项只与**消息条数**有关,与报文大小、消息内容都无关,因此在**小报文 × 高吞吐流式**(日志 / 遥测 / 批量采集)场景里由固定开销主导,构成瓶颈。 实测(128 B 报文、每流连发 128 条、并发 32): - **strace 计数 2.10 次** **`write(2)`/消息**(14516 ÷ 6912;多出的约 0.1 次来自 HEADERS / WINDOW_UPDATE 等控制帧,该组按 `-concurrency 1` 采集以便反推消息总数)。若写路径存在 32 KiB 级别的聚合,128 B 消息要攒 256 条才满一个水位,同样的消息量也只需约 **1/128** 次 `write(2)`(一个 32 KiB 批会被 http2 切成两帧、每帧各 flush 一次),低两个数量级,说明每条消息各自触发两次独立写入; - **端到端压测三个特征同向**:并发从 32 提到 128(4×)QPS 不增反降、消息吞吐随报文增大单调下降、每流连发条数涨 4 倍吞吐也不涨,都指向吞吐上限由**每条消息的固定开销**决定,而不是由并发或连发条数决定; - **block profile(同一 20 s 窗口差分)显示阻塞时间 99.77 % 落在三个等待上**,几乎没有 on-CPU 计算,只能说明耗时是等待而不是编解码计算= 建议引入一个**默认关闭**的可选批缓冲 `streamBufferWriter`(32 KiB 水位):把连续小消息攒进一块缓冲,按水位或 `CloseRequest` 冲刷,把交接次数与 syscall 次数摊薄到水位级别,以 `WithWriteBuffering()` 显式 opt-in。 ## 相关代码 | 位置 | 问题 | | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------------------- | | [`protocol_grpc.go#L299-L300`](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/protocol_grpc.go#L299-L300)(`grpcClient.NewConn`) | 未开启 `WriteBuffering` 时 `streamWriter = call`,协议层与 `io.Pipe` 之间没有缓冲层,每次 `Send` 直接下沉 | | [`envelope.go#L147-L161`](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/envelope.go#L147-L161)(`envelopeWriter.write`) | 一条消息拆成 prefix 与 body 两次 `Write`,两次都立即下沉,是 gRPC wire 每消息 2 次交接的来源 | | [`duplex_http_call.go#L108-L126`](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/duplex_http_call.go#L108-L126)(`duplexHTTPCall.Write`) | 每次直写 `io.PipeWriter`,中间没有任何攒批逻辑 | | [`protocol_triple.go#L587-L595`](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/protocol_triple.go#L587-L595)(`tripleUnaryMarshaler.write`) | Triple 编码流式只有一次 body 写(客户端流式也走这个 marshaler),同样每条消息各付一次交接,没有聚合 | | `golang.org/x/[email protected]/http2/transport.go` L1566-L1576(`writeRequestBody`) | 每个 DATA 帧立即 `cc.bw.Flush()`,一次管道交接对应一次真实 `write(2)` | ## 建议优化 在 `duplexHTTPCall` 之上插一层缓冲 `streamBufferWriter`,挂在**未走 fast path 的分支**(流式 + 非 fast path 的 unary),默认关闭。 ### 1. 聚合层 ```go const defaultStreamWriteBufSize = 32 << 10 // 32 KiB type streamBufferWriter struct { mu sync.Mutex next io.Writer // wrapped duplexHTTPCall limit int buf *bytes.Buffer err error closed bool } ``` - **`Write`**:先把 `p` 追加进 `buf`;一旦 `buf.Len() >= limit`,在**同一临界区内**调 `flushLocked` 把整批一次下沉; - **单条 ≥ 水位**:先把缓冲里已有的小消息刷出去,再把这条大消息**直写**底层,避免一条大消息在内存里暂存两次; - **`Flush`**:空缓冲直接返回,不在 `io.Pipe` 上白跑一次交接,也不给 http2 制造空帧的机会; - **`flushLocked`** **失败**:**首错即冻结**,记下 `err` 并置 `closed`,之后所有 `Write` / `Flush` / `Close` 一致返回同一个错误; - **`Close`**:在**同一临界区**内完成"刷尽 + 置 closed",封口之后的 `Write` 返回 `io.ErrClosedPipe`。 ### 2. 装配 两条协议客户端的流式分支都接上,默认不动([protocol\_grpc.go#L308-L325](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/protocol_grpc.go#L308-L325)、[protocol\_triple.go#L278-L285](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/protocol_triple.go#L278-L285)): ```go call = newDuplexHTTPCall(ctx, g.HTTPClient, g.URL, spec, header) streamWriter = call if g.WriteBuffering { writeBuffer = newStreamBufferWriter(call) streamWriter = writeBuffer } ``` - `writeBuffer` 字段与 `envelopeWriter.writer` 由**同一个** **`streamWriter`** **变量派生**,天然成对,不会出现"writer 指向缓冲而字段为 nil,`CloseRequest` 不知道该刷缓冲、尾包被静默丢掉"的错配; - `CloseRequest` 只在 `writeBuffer != nil` 时先 `Close()` 刷尽,随后**无论 flush 成败都执行** **`CloseWrite`**,两个错误取第一个非 nil(flush 优先); - 未开启时 `writeBuffer = nil`、`marshaler.writer = call`,与现状**逐字节等价**; - unary fast path、服务端路径一律不改。unary 未走 fast path 时会被顺带套上缓冲,因为只有一条消息,聚合最多省 1 次交接而要多付一次拷贝与首包延迟,收益接近零,已在 option 的文档注释中如实说明。 ### 3. 方案边界与代价 用**首包延迟、一次内存拷贝、一把锁**,换掉每条消息两次同步交接与两次 syscall。因此: - **适用**:流式客户端(gRPC wire 默认路径 + Triple 编码)、单流连发大量小消息、吞吐优先于即发即达; - **不适用 / 不生效**:unary fast path(走 `newUnaryFastPathCall`,不经过 `io.Pipe`)、服务端路径、单条 ≥ 32 KiB 的报文(走直写路径)、低首包延迟敏感场景; - **多 goroutine 并发写同一条流**:聚合后每次写都要先取锁,`mu` 会成为新的争用点,需要自行权衡; - **默认关闭**:未开启用户零感知。 ## 测试与设计约束 **设计约束:** 1. **关闭必须在同一临界区内刷尽并封口**:`Close()` 在持锁状态下完成刷尽 + 置 closed,于是存在一个可断言的不变量:**到达底层的字节数 == 被接收的 `Write` 数 × 消息长度**。把关闭拆成两段(先 Flush 再 CloseWrite)会让这个不变量失效,而且失效方式是静默的,这正是它必须原子化的原因; 2. **直写之前先刷缓冲**:保序是聚合层自己的责任。把大消息直写当成一条捷径时,很容易忽略捷径不能插队,所以判断顺序必须是先处理缓冲、再处理直写; 3. **首次错误即冻结状态**:流是一次性资源,写失败后任何再试一次都没有语义,只会把故障暴露时间推迟。冻结的另一个好处是错误**可复现**:第一次失败之后的每一次调用都返回同一个错误; 4. **错误类型与错误来源逐字对齐**:聚合层加进来之后,上层不应该观察到任何新形态的错误。所以封口后写沿用 `io.ErrClosedPipe`,管道真正关闭时仍由 `duplexHTTPCall.Write` 统一转成 `io.EOF`,聚合层不新增、不改写任何错误; 5. **空操作不下沉**:空缓冲 flush 直接返回,既不在 `io.Pipe` 上白跑一次交接,也不给 http2 制造空帧的机会。这条同时让 `Flush` 的重复调用天然幂等; 6. **兜底责任不转移给聚合层**:`CloseWrite` 里的 `ensureRequestMade()` 继续负责零消息流也要把请求发出去,聚合层不接管这个职责,也就不可能破坏它。 这些约束现有测试已守住了大部分,还需要补下面 13 个**单元测试**: | 要补的测试 | 功能 | | --------------------------------------------------------- | --- | | `TestStreamBufferWriterStickyErrorOnEveryExit` | 首错冻结之后 `Write` / `Flush` / `Close` 三个出口返回同一个错误 | | `TestStreamBufferWriterWriteReportsFlushFailure` | 水位 flush 失败但本次字节已入缓冲,返回 `(len(p), err)` 而不是 nil | | `TestStreamBufferWriterLargeMessageErrorPropagates` | 大消息直写失败时,上层拿到的是同一个错误,而不是 nil | | `TestStreamBufferWriterCloseFlushFailureReturnsError` | `Close` 尾刷失败时返回该错误,并置 `closed`、不重复尝试 | | `TestCloseRequestClosesWriteSideAfterFlushFailure` | flush 失败时 `CloseWrite` 仍被调用,对端能读到流结束 | | `TestStreamBufferWriterConcurrentCloseAndFlush` | `Close` 与 `Flush` 并发时同一批数据不会被发两次 | | `TestStreamBufferWriterConcurrentBufferBound` | 并发 `Send` 下缓冲以 `limit` 为界,不会无界增长 | | `TestStreamBufferWriterLargeMessageFlushesPendingFirst` | 大消息直写之前先把缓冲里的小消息刷出去,不插队 | | `TestStreamBufferWriterShortWriteIsFailure` | 下游短写(`n < len(p)` 且 err == nil)按失败处理并冻结,不静默丢字节 | | `TestStreamBufferWriterCloseAfterCloseFailure` | flush 失败之后再次 `Close` 返回同一个粘滞错误 | | `TestStreamBufferWriterFlushAfterClose` | `Close` 成功之后再 `Flush` 返回 nil,不报 `io.ErrClosedPipe` | | `TestWriteBufferingCoversUnaryNonFastPath` | unary 未走 fast path 时同样被套上缓冲 | | `TestWriteBufferingEmptyStreamStillSendsRequest` | 一条消息都不发就直接 `CloseRequest`,服务端仍能收到这次请求 | ## 复现方式 压测端到端 AB(真实 dubbo-go 服务端,AB 只切 `-buffering`): ```bash # 服务端 go run ./tools/benchmark/server/dubbo-go --serialization protobuf --compression none --port 20123 # 客户端:-buffering 关闭 = 优化前,打开 = 优化后 go run ./tools/benchmark/stream_ab -addr 127.0.0.1:20123 \ -payload 128 -msgs 128 -concurrency 32 -duration 60s -warmup 10s # 单流自检(先确认链路通) go run ./tools/benchmark/stream_ab -addr 127.0.0.1:20123 -probe ``` 协议层 micro-bench(`os.Pipe` 真 fd,同报文档位、同总字节,只换写策略): ```bash cd protocol/triple/triple_protocol go test -run '^$' -bench 'BenchmarkStreamWrite' -benchmem -count=10 benchstat benchstat_input_before.txt benchstat_input_after.txt ``` syscall 计数(真实链路,单并发以便用成功流数 × 每流条数反推消息总数): ```bash strace -c -f -e trace=write -o strace.txt \ go run ./tools/benchmark/stream_ab -addr 127.0.0.1:20123 \ -payload 128 -msgs 128 -concurrency 1 -duration 5s -warmup 0s ``` > `-f` 不可省略:Go 会把 goroutine 调度到多个线程上,不跟随线程会漏掉大部分写。 profiling(CPU 与 block 需先对齐时间基准,block 侧取同窗口前后两次快照做差分): ```bash go run ./tools/benchmark/stream_ab -addr 127.0.0.1:20123 \ -payload 128 -msgs 128 -concurrency 32 -warmup 10s -duration 120s -pprof 127.0.0.1:6060 curl -s -o cpu.prof "http://127.0.0.1:6060/debug/pprof/profile?seconds=20" # block 侧单独一轮,-blockrate 1 打开采样 curl -s -o block_t0.prof "http://127.0.0.1:6060/debug/pprof/block" sleep 20 curl -s -o block_t1.prof "http://127.0.0.1:6060/debug/pprof/block" go tool pprof -proto -base block_t0.prof -output block_window.prof block_t1.prof go tool pprof -http=:8081 -sample_index=delay block_window.prof ``` ## 参数来源 写缓冲这个形态在主流 RPC / 消息中间件里是标准做法,水位(32 KiB)与主流一致;本方案默认关闭比 gRPC Go / Kitex(默认 32 KB 开启)更保守,其他参数参考了业界较为合理的数值。 | 项目 | 对应开关 | 水位与默认值 | | -------------------------------- | ------------------------------------------------------ | ----------------------------------- | | gRPC Go | `WithWriteBufferSize()` | 默认 32 KB,传 `0` 关闭 | | gRPC C++ | `GRPC_ARG_HTTP2_WRITE_BUFFER_SIZE` / `set_buffer_hint` | 流式写缓冲为 opt-in | | Kitex | `WithGRPCWriteBufferSize()` | 默认 32 KB | | tRPC-Go tnet websocket | Combined Writes Optimization | 小消息合并为单次 syscall | | Kafka producer | `batch.size` / `linger.ms` | `batch.size` 默认 16384 | | OpenTelemetry BatchSpanProcessor | `max_export_batch_size` | 每批 512 条,遥测 span 真实大小 200 B \~ 1 KB | ## 参考链接 - gRPC-Go `WithWriteBufferSize`:<https://pkg.go.dev/google.golang.org/grpc> - gRPC C++ Performance Notes(streaming write buffering):<http://grpc.github.io/grpc/cpp/md_doc_cpp_perf_notes.html> - gRPC-Go performance improvements(flush syscall for every data frame):<https://grpc.io/blog/grpc-go-perf-improvements/> - Kitex client options:<https://www.cloudwego.io/docs/kitex/tutorials/options/client_options/> - tRPC-Go tnet websocket(Combined Writes):<https://pkg.go.dev/trpc.group/trpc-go/tnet/extensions/websocket> - Kafka Producer Configs(`batch.size` / `linger.ms`):<https://kafka.apache.org/41/configuration/producer-configs/> - OpenTelemetry BatchSpanProcessor 默认值:<https://opentelemetry.io/docs/zero-code/dotnet/configuration/> -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
