GitHub user lizining1231 edited a comment on the discussion: [OSPP 2026] Triple 协议性能分析与优化
# 性能瓶颈分析报告(三)Triple 流式客户端写路径批缓冲专项 ## 目录 一、摘要 二、瓶颈代码定位 三、多工具交叉验证瓶颈 四、方案设计与演进 五、方案验证与实施 六、预估收益与方案边界 七、测试构造情况说明 八、参考链接 ## 一、摘要 Triple 是 Dubbo-go 的默认协议层,其流式客户端(`CallBidiStream` / `CallClientStream`)支持"一条流上连续发送多条消息"的用法。这类用法的共性是:**消息条数多、连发节奏快**,单条消息则不一定大。 而当前实现里,每次 `Send` 都直接把消息写进 `duplexHTTPCall` 底层的 `io.Pipe`:它是无缓冲的跨 goroutine 同步交接,写端必须等读端消费完才返回,因此**每发一条消息就要做两次同步握手(prefix 与 body 各一次)**。消息越多,握手次数越多,per-message 代价与消息条数成正比,在高吞吐连发时成为隐藏的固定开销,构成**小报文 × 高吞吐流式**场景下的潜在瓶颈。 本文针对 Triple 流式客户端发送路径的**结构性开销**,验证瓶颈的确存在,确认**每条消息两次无缓冲同步交接、并落成两次独立 socket 写**是主因;在此基础上提出一个**默认关闭**的可选批缓冲 `streamBufferWriter`(32 KiB 水位)方案:把连续小消息聚合、按水位 / `CloseRequest` 冲刷,摊薄交接次数与 syscall 次数。方案以 `WithWriteBuffering()` option 显式开启,默认关闭。 ## 二、瓶颈代码定位 前两期的优化都是从**相较 gRPC 的基准差距**出发反推瓶颈,本期路径不同:triple unary 落地后,unary 侧的固定开销已被吸收,剩余可优化范围收敛到**流式客户端**,而流式路径并没有现成的数据可供入手分析。因此本期先从静态源码入手、再设计测试验证瓶颈是否真实存在。 ### 2.1 流式 Send 的写路径没有聚合 `CallBidiStream`([client.go L186-L195](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/client.go#L186-L195))是"一条流上连续 `Send` N 条消息"的流式用法。顺着源码往下读四个环节,就能看清一次 `Send` 到底做了什么。 **第一,连接装配**:`grpcClient.NewConn` 建流式连接时,`streamWriter` 直接就是 `duplexHTTPCall`,协议层与管道之间没有任何缓冲层(默认配置下,[protocol\_grpc.go L299-L325](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/protocol_grpc.go#L299-L325)): ```go call = newDuplexHTTPCall(ctx, g.HTTPClient, g.URL, spec, header) streamWriter = call // 协议层的 writer 就是 duplexHTTPCall ... conn := &grpcClientConn{ call: call, marshaler: grpcMarshaler{ envelopeWriter: envelopeWriter{ writer: streamWriter, // envelope 的每次 Write 都直接下沉 ... }, }, } ``` **第二,Send 入口**:`Send` 只做一件事,把消息交给 marshaler([protocol\_grpc.go L370-L375](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/protocol_grpc.go#L370-L375)): ```go func (cc *grpcClientConn) Send(msg any) error { if err := cc.marshaler.Marshal(msg); err != nil { return err } return nil } ``` **第三,协议层组帧**:marshaler 把消息编好之后,**一条消息要写两次**,先写 5 字节 prefix(flag + 消息长度),再用 `io.Copy` 写 body([envelope.go L147-L161](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/envelope.go#L147-L161)): ```go func (w *envelopeWriter) write(env *envelope) *Error { prefix := [5]byte{} prefix[0] = env.Flags binary.BigEndian.PutUint32(prefix[1:5], uint32(env.Data.Len())) if _, err := w.writer.Write(prefix[:]); err != nil { // 第 1 次:prefix ... } if _, err := io.Copy(w.writer, env.Data); err != nil { // 第 2 次:body ... } return nil } ``` **第四,下沉写**:这两次 `Write` 落到 `duplexHTTPCall`,每次都是直写 `io.PipeWriter`,中间没有任何攒批逻辑([duplex\_http\_call.go L108-L126](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/duplex_http_call.go#L108-L126)): ```go func (d *duplexHTTPCall) Write(data []byte) (int, error) { d.ensureRequestMade() // 首次调用时触发后台 makeRequest goroutine ... bytesWritten, err := d.requestBodyWriter.Write(data) // 直写 io.PipeWriter ... } ``` 总结来看,四个环节各有开销,且都是每消息的开销。 | 环节 | 代码行为 | 每消息开销 | | ---------------------- | ----------------------------------------------- | ----------------------- | | `Send(msg)` | `grpcClientConn.Send` → `grpcMarshaler.Marshal` | 1 次 marshal | | 协议层组帧 | `envelopeWriter.write` 写 prefix + body | **2 次下沉 `Write`** | | `duplexHTTPCall.Write` | 每次直写 `io.PipeWriter` | **2 次跨 goroutine 同步交接** | | 传输层落 socket | 每个 DATA 帧立即 flush(3.3 实测) | **2 次 `write(2)`** | **初步结论**:这段调用链里没有任何聚合点,一个 `Send` 的固定成本结构是明确的,1 次 marshal、2 次下沉写、2 次同步交接、2 次 socket 写,而后三项只与**消息条数**有关,与消息内容、报文大小都无关。所以可以大致判断:在小报文 × 高吞吐流式这类条数多、单条小、连发快的场景里,每条消息都要摊一次这段固定开销,写路径可能存在瓶颈。这个判断来自静态读码,是否真的构成瓶颈、开销落在哪一层,还需要在真实链路上实测确认。 ### 2.2 从 Send 到 socket 的完整时序 把 2.1 的四段代码摊平成一条流水线(每一步都对应上面引用的源码)。一次 `Send` 会拆成 prefix 与 body 两次 `Write`,**两次都是串行的:① 一路执行(含落 socket)之后才轮到 ②**,而这两条路径的形状完全一样。`wrCh` / `rdCh` 是 `io.Pipe` 内部的两个 channel: ```mermaid flowchart LR S["Send(msg)<br/>marshal 得到 body"] --> A["① prefix<br/>5 B"] A --> A1["duplexHTTPCall.Write<br/>直写 io.PipeWriter,写端阻塞"] A1 --> A2["writeRequestBody 取走这一段<br/>立即 Flush"] A2 --> A3["socket write(2)"] A3 --> B["② body<br/>N B<br/>(① 落完 socket 才轮到)"] B --> B1["duplexHTTPCall.Write<br/>直写 io.PipeWriter,写端阻塞"] B1 --> B2["writeRequestBody 取走这一段<br/>立即 Flush"] B2 --> B3["socket write(2)"] ``` 图上可以看出:**一条消息要走两次完整的管道交接、落成两次 socket 写**,而且这两次是串行的,没有任何合并的机会。把这条流水线拆开,一次 `Send` 会在 ① ② 两条路径上各完整走一遍下面六步(下面以默认的 gRPC wire 为例,它有 5 字节 prefix;Triple 编码没有 prefix,只走 ② 那一条): 1. **序列化**:`Send` 调 `grpcMarshaler.Marshal` 把消息编成 `envelope`(body),拿到待写的字节,此刻还没有碰到管道; 2. **组帧**:`envelopeWriter.write` 用同一个 `writer` 发起第一次 `Write`,写 5 字节 prefix(flag + 大端长度); 3. **下沉到调用对象**:这次 `Write` 落到 `duplexHTTPCall.Write`,它先 `ensureRequestMade()` 拉起后台 `makeRequest` goroutine(首次调用时),校验 ctx 未取消,再用 `d.requestBodyWriter` 写数据,而这个 `requestBodyWriter` 就是 `io.PipeWriter`; 4. **写端 park(第一个阻塞点)**:`io.Pipe` 不做缓冲,写端把这段切片交给读端后,立刻等读端回传已取走的字节数,写端在这两处 channel 操作上阻塞,直到读端出现; 5. **读端取走并落 socket**:读端不是业务代码,是 `net/http` 从 `http.Request.Body` 读数据的 `writeRequestBody`。它把这段拷进自己的缓冲区、回传字节数,再按 HTTP/2 帧编码,每个 DATA 帧立即 `Flush`,最终经 `net.Conn.Write` 落成 `write(2)`; 6. **写端返回,发起第二次 `Write`**:第 4 步返回之后 `duplexHTTPCall.Write` 才返回,`envelopeWriter.write` 随即用 `io.Copy` 写 body,把第 3 到第 5 步整个重做一遍。 - 写端 park 的那段时间是记在调度器头上的:pprof 采样会落在 `runtime.selectgo` 而不是 `io.(*pipe).write`; - 传输层对每个 DATA 帧立即 flush,所以一次管道交接最终对应一次 `write(2)`,实测约 2.10 次/消息。 `io.Pipe` 的写入不是把数据拷进缓冲区就返回,而是**与读端做一次同步交接**。标准库的实现只有两个 channel:`wrCh`(传递数据切片)与 `rdCh`(回传已取走的字节数),写之前还要取 `wrMu` 把同一个管道上的 `Write` 串行化([pipe.go L76-L96](file:///home/lizining/go-sdk/go/src/io/pipe.go#L76-L96)): ```go func (p *pipe) write(b []byte) (n int, err error) { select { case <-p.done: return 0, p.writeCloseError() default: p.wrMu.Lock() // 串行化同一管道上的写 defer p.wrMu.Unlock() } for once := true; once || len(b) > 0; once = false { select { case p.wrCh <- b: // ① 把这段数据交给读端,写端在此阻塞 nw := <-p.rdCh // ② 等读端回传已取走的字节数,写端再次阻塞 b = b[nw:] n += nw case <-p.done: return n, p.writeCloseError() } } return n, nil } ``` 读端与之对应:从 `wrCh` 取值、`copy` 进自己的缓冲区、再把实际取走的字节数发到 `rdCh`([pipe.go L50-L65](file:///home/lizining/go-sdk/go/src/io/pipe.go#L50-L65))。由此有两点结论: - 一次 `Write` 至少要**两次跨 goroutine 唤醒**(读端一次、写端一次),读写两端交替推进,写端必须等读端把这段取走才能返回; - 交接以**切片**为单位:读端缓冲区一次装不下时,写端会在循环里把剩余部分再交接一次(`b = b[nw:]`)。真实链路的读端是 http2 的 `writeRequestBody`,它用约 16 KiB 的 scratch buffer 读取,因此 5 字节 prefix 与 128 B body 都是整段一次交接完。 这条链路上的读端不是业务代码,而是 `net/http`:`newDuplexHTTPCall` 把 `io.Pipe` 的读端直接挂在 `http.Request.Body` 上([duplex\_http\_call.go L76-L94](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/duplex_http_call.go#L76-L94)),由 h2c 传输层按需读取、交给 HTTP/2 帧编码器。所以每次 `Send` 都要等传输层把这一段读走才能返回,**同步交接次数随消息条数线性增长**。 ### 2.3 对瓶颈的预期 读源码只能说明代码长这样,还不能说明它一定构成瓶颈。把 2.1、2.2 的代码结论翻译成能在真实链路上检验的预期瓶颈,就是下面内容,测试也是参考它们设计的。 1. **管道不做合并**:`io.Pipe` 的职责只是把这一段数据交给读端,它不攒批,所以写一次就是交接一次。这一条已由 2.2 的源码坐实; 2. **每条消息写两次 socket**:协议层把一条消息拆成 prefix 与 body 两次 `Write`,传输层对每个 DATA 帧又立即 flush,因此 `write(2)` 的次数应当跟着**消息条数**走,而不是跟着**字节数**走。如果写路径中间存在攒批,写次数就会与消息条数脱钩、远少于消息条数; ## 三、多工具交叉验证瓶颈 ### 3.1 端到端压测:现状吞吐与延迟表现 客户端以单**流连发 N 条消息 → `CloseRequest` → 收全 N 条响应**为一个操作单元计一次 QPS;表中的**消息吞吐 = QPS × 每流连发条数**,由两者换算得到,衡量单位时间内真正处理的消息条数。 **矩阵**:5 种报文(128 B / 256 B / 512 B / 1 KiB / 4 KiB)× 3 档并发(32 / 64 / 128)× 2 档连发(128 / 512)= 30 组数据,覆盖多种参数区间。 #### (1)连发 128 条 × 5 报文 × 3 并发 **128 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 200.6 | 25 677 | 222.59 | | 64 | 185.8 | 23 782 | 511.83 | | 128 | 165.7 | 21 210 | 1125.09 | **256 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 186.8 | 23 910 | 259.65 | | 64 | 143.8 | 18 406 | 716.16 | | 128 | 139.8 | 17 894 | 1382.00 | **512 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 138.9 | 17 779 | 324.66 | | 64 | 127.6 | 16 333 | 804.96 | | 128 | 124.0 | 15 872 | 1598.27 | **1024 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 102.9 | 13 171 | 433.95 | | 64 | 99.8 | 12 774 | 937.74 | | 128 | 94.5 | 12 096 | 1911.53 | **4 KiB** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 63.3 | 8 102 | 660.56 | | 64 | 58.6 | 7 501 | 1383.48 | | 128 | 57.9 | 7 411 | 2858.83 | #### (2)连发 512 条 × 5 报文 × 3 并发 **128 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 46.3 | 23 711 | 954.48 | | 64 | 44.1 | 22 564 | 1949.47 | | 128 | 41.9 | 21 432 | 4150.53 | **256 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 41.7 | 21 345 | 1066.70 | | 64 | 39.2 | 20 076 | 2249.60 | | 128 | 38.6 | 19 758 | 4523.50 | **512 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 35.4 | 18 125 | 1213.60 | | 64 | 35.1 | 17 971 | 2488.90 | | 128 | 34.1 | 17 459 | 5158.20 | **1024 bytes** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 28.0 | 14 336 | 1447.00 | | 64 | 23.8 | 12 186 | 3341.80 | | 128 | 25.8 | 13 210 | 6054.30 | **4 KiB** | Concurrency | QPS (流/s) | 消息吞吐 (msgs/s) | P99 (ms) | | ----------- | --------: | ------------: | -------: | | 32 | 16.4 | 8 417 | 2286.95 | | 64 | 18.4 | 9 441 | 4275.64 | | 128 | 15.3 | 7 859 | 9969.09 | #### 三个显著特征 **特征一:并发提高 → QPS 不增反降(系统近饱和)** 以 128 B / 连发 128 为例,并发从 32 提到 128(4×),QPS 反而从 200.6 降到 165.7,消息吞吐从 25 677 降到 21 210 msgs/s;4 KiB / 连发 512 同样如此(16.4 → 15.3 流/s)。新增并发没有带来等比收益,反而因为每条流都在做同样的同步交接、争抢同一份写路径资源,单流效率随之下降,这说明现状已经处在某种饱和状态。 **特征二:消息吞吐随报文增大单调下降** 并发 32 下,连发 128 条的消息吞吐从 25 677(128 B)降到 8 102(4 KiB),下降约 3.2×。但把它换算成字节吞吐,128 B 约 3.3 MB/s,4 KiB 约 33 MB/s,**反而涨了约 10 倍**。若瓶颈是搬运字节的带宽,字节吞吐应当先封顶,而它随报文增大一路上升,说明真正被放大的是**每条消息摊到的固定成本**:小消息每字节分摊的固定开销高、大消息分摊低,于是消息吞吐随报文增大单调下降。 **特征三:增加连发条数 → 消息吞吐不再增长** 把每流连发条数从 128 提到 512(4×),各组合的**消息吞吐基本持平**:128 B / c32 从 25 677 微降到 23 711 msgs/s(-7.7 %),512 B / c64 从 16 333 升到 17 971 msgs/s(+10.0 %),4 KiB / c64 从 7 501 升到 9 441 msgs/s。多连发 4 倍的消息,单位时间能处理的消息数并没有变多,**吞吐上限由每条消息的固定开销决定,与流内连发多少条无关**。 延迟侧的表现更直接:128 B / c32 下单流 P99 从 223 ms 涨到 954 ms(4.3×),与连发条数 4× 几乎等比。若写路径存在聚合,512 条消息的写次数应被摊薄到远小于 4×,单流耗时不会线性放大。 这三个特征互相印证现状写路径存在的瓶颈:每条消息都要走两次独立的同步交接,吞吐上限由**固定开销**决定而不是由并发或连发条数决定。下文用 CPU profile 与 strace 计数进一步确认这笔开销的来源。 ### 3.2 客户端 CPU profile 与 block profile 客户端 `stream_ab` 以 128 B 报文、每流连发 128 条、并发 32 持续压测(未开启任何写缓冲),同时通过 `net/http/pprof` 抓取客户端进程的 CPU profile: <img width="1917" height="910" alt="image" src="https://github.com/user-attachments/assets/35d68a69-4c39-482f-a34b-5cd3eb021375" /> **解读**: - **写路径链路**:`http2.(*clientStream).writeRequestBody` cum 20.31 %,其下 `bufio.(*Writer).Flush` cum 14.50 % → `http2.writeWithByteTimeout` cum 13.69 % → `net.(*conn).Write` cum 13.29 %,最终落到 `internal/runtime/syscall.Syscall6`(该函数为全进程 flat 23.27 %,同时承担读侧 syscall)。整条栈就是把一段数据写进 socket 这一件事,中间没有任何应用层逻辑; - **并列的宽条**是 `internal/runtime/syscall.Syscall6`(flat 23.27 %)与 `runtime.futex`(flat 17.76 %),合计 **41.03 %**:前者是进内核写 socket,后者是 goroutine 的唤醒与调度,火焰图右侧 `runtime.mcall → park_m → schedule → findRunnable → stealWork` 这条栈清晰可见; - **确实存在管道交接**:`io.(*pipe).read → runtime.selectgo` 这条栈对应传输层从管道取数据的那一侧,说明写路径确实是写端写管道 → 读端取走 → 落 socket 一段一段推进的; - **编解码很小不是主要成本**:收侧 `grpcUnmarshaler.Unmarshal` cum 20.31 %,写侧 `envelopeWriter.Marshal` 的占比远小于它下面的 syscall 部分。 **但是这里不能只看 CPU profile**,CPU profile 只统计 on-CPU 时间,`io.Pipe` 交接里的等待属于 off-CPU,火焰图上天然看不到。 block profile <img width="1916" height="907" alt="image" src="https://github.com/user-attachments/assets/db7b3e20-f552-44a1-8183-97e18595ebda" /> **占比最高的三个帧都不是业务帧,而是三个等待**:`runtime.selectgo` **62.90 %**、`runtime.chanrecv1` **21.57 %**、`sync.(*Cond).Wait` **15.30 %**,三者合计 **99.77 %**;而写路径自身的帧 `io.(*pipe).write` 只有 **0.51 %**(8.24 s),它的读端 `io.(*pipe).read` 只有 **0.06 %**(0.98 s)。它排除了别的优化方向, - 如果这里测出的是**计算占满**(编解码、拷贝在吃时间),方案应偏向于优化 codec、减少拷贝; - 实测是**等待占满**(99.77 %),说明这条路径的耗时形态是等待形态而不是计算形态,不是在编解码上算不过来,而是一条消息一次往返地等过来。方案方向因此定为**只能改交接结构,把次数降下来**。 指令 ```bash # 轮 1:采 CPU stream-ab -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" # 轮 2:同参数单独一轮采 block stream-ab -payload 128 -msgs 128 -concurrency 32 -warmup 10s -duration 120s \ -pprof 127.0.0.1:6060 -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" # 差分出与 CPU 同窗口的 block profile 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 ``` ### 3.3 strace write 计数 **采集命令**(真实链路,单并发以便用成功流数 × 每流条数反推出确定的消息总数): ```bash strace -c -f -e trace=write -o strace.txt \ stream-ab -payload 128 -msgs 128 -concurrency 1 -duration 5s -warmup 0s ``` > 其中 `-f` 不可省略:Go 会把 goroutine 调度到多个线程上,不跟随线程会漏掉大部分写 **采集结果**:客户端报告 `Success = 54` 条流,即共发送 **6912 条** 128 B 消息;strace 统计 `write` 调用 **14516 次**: ``` % time seconds usecs/call calls errors syscall ------ ----------- ----------- --------- --------- ---------------- 100.00 1.694243 116 14516 write ``` **折算:2.10 次 `write(2)` / 消息**(14516 ÷ 6912)。这个数字与源码是相符的: | 环节 | 代码依据 | 每条消息代价 | | ----------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------- | | 协议层组帧 | `envelopeWriter.write` 先写 5 字节 prefix(flag + 长度),再用 `io.Copy` 写 body([envelope.go L147-L161](https://github.com/apache/dubbo-go/blob/develop/protocol/triple/triple_protocol/envelope.go#L147-L161)) | **2 次 `Write` 进管道** | | 传输层落 socket | x/net/http2 `writeRequestBody` 对每个 DATA 帧都立即 `cc.bw.Flush()`(注释原文 “TODO(bradfitz): this flush is for latency, not bandwidth”),每次 Write 对应一次真实 `write(2)`(`golang.org/x/[email protected]/http2/transport.go` L1566-L1576) | **2 次 `write(2)`** | 实测 2.10 略高于 2,多出的约 0.1 次/消息来自 HEADERS、WINDOW\_UPDATE 等控制帧(54 条流合计 692 次)。 **解读**: - **次数**:`write` 调用数与消息条数成正比(2.10 次/条),而不是 1/256 次/条。若写路径存在 32 KiB 级别的聚合,128 B 消息要攒 256 条才满一个水位,同样的消息量只需约 **1/256** 的 `write(2)`,**差两个数量级**,说明每条消息各自触发两次独立写,没有任何合并; - **耗时需谨慎推理**:`usecs/call = 116 µs` 是在 ptrace 跟踪下测得的,每次 syscall 进入 / 退出都要与 tracer 交互,绝对值被明显放大。它只能用来说明**每条消息都要完整走一遍 syscall 路径**,不能当作真实业务耗时。 ### 3.4 交叉验证小结 - **端到端压测**给出外部表现:并发提高不能提高吞吐(特征一)、报文越大吞吐越低(特征二)、连发条数涨 4 倍吞吐也不涨(特征三),三条都指向**吞吐上限由每条消息的固定开销决定**; - **CPU profile**体现了时间去向:`http2.(*clientStream).writeRequestBody` cum 20.31 %;`internal/runtime/syscall.Syscall6` flat 23.27 % 与 `runtime.futex` flat 17.76 % 合计 **41.03 %**,代价集中在**进内核写 socket**与**goroutine 唤醒调度**,而不是应用逻辑; - **strace 计数**给出次数:2.10 次 `write(2)`/消息,与 prefix + body 两次 `Write`、每个 DATA 帧立即 flush 的代码路径逐层对应,与 32 KiB 聚合应有的量级差两个数量级; - **block profile** 给出阻塞构成:同一个 20 s 窗口的差分口径下,流式调用的阻塞时间 **99.77 %** 落在三个等待原语上(`runtime.selectgo` 62.90 %、`runtime.chanrecv1` 21.57 %、`sync.(*Cond).Wait` 15.30 %),几乎没有 on-CPU 的计算;而一次 `Write` 对应读端一段、管道不做合并这一交接语义直接来自标准库源码([pipe.go L76-L96](file:///home/lizining/go-sdk/go/src/io/pipe.go#L76-L96)),不需要额外的合成路径实验来佐证。该占比会因写端先 park 还是读端先 park 在不同窗口之间互补翻转(实测 `io.(*pipe).write` 0.51 % \~ 5.89 %),**故仅作定性参考,不作为瓶颈的定量证据**。 把这四者串起来,瓶颈是**结构性**的:流式客户端写路径上,每条消息都要做两次无缓冲的跨 goroutine 同步交接,并因此落成两次独立的 socket 写,且没有任何聚合机制来摊薄它。这也决定了后续方案的形状,在不改变 on-wire 语义的前提下,把多次小写聚合成一次大写。 ## 四、方案设计与演进 ### 4.1 问题定位 瓶颈锁定在流式客户端每次 `Send` 的 `io.Pipe` 同步握手(per-message 代价 ∝ 消息数),无任何聚合。 ### 4.2 设计目标 - **默认零影响**:`WithWriteBuffering()` 默认关闭,未开启时 `writeBuffer=nil`、`marshaler.writer=call`,与现状**逐字节等价**; - **语义不变**:On-wire 帧格式不变(默认 gRPC 编码仍是 flag + 长度的 5 字节 prefix,Triple 编码仍是裸 body),只是把多次小 `Write` 合并成一次大 `Write`;不改公开接口签名、不改 `Codec` 接口、无 on-wire 格式变更。需要说清的是**帧格式不变不等于到达时序不变**:批边界会改变对端开始消费这条流的数据的时刻(小流甚至要等到 `CloseRequest` 才发出第一个字节),这是本方案主动付出的代价之一; - **不丢消息**:`CloseRequest` 必须刷尽剩余缓冲,尾包完整送达(并发契约守护); - **opt-in 哲学对齐行业**:gRPC-Go / Kitex 均为可关闭的写缓冲开关。 ### 4.3 并发相关的边界条件 `StreamingClientConn` 接口契约要求 `Send` / `RequestHeader` / `CloseRequest` **可互相并发**,因此聚合层的关闭语义必须与关闭动作在同一临界区内封口。如果把 `CloseRequest` 写成先 `Flush()` 再 `CloseWrite()`(两次独立加锁),就会出现**静默丢消息**窗口:一个 `Send` 已被接收、数据还留在缓冲里,紧接着 `CloseWrite` 关闭管道 → 消息凭空消失;而在无缓冲的现状下,那个并发 `Send` 会拿到 `io.ErrClosedPipe` **报错**。聚合层不能把报错变成静默丢失。 设计约束由此确定:`Close()` 在**同一临界区**内完成刷尽 + 置 closed,封口之后的 `Write` 返回 `io.ErrClosedPipe`(错误类型与现状逐字一致);`CloseRequest` 无论 flush 成败都执行 `CloseWrite`,错误取第一个非 nil。 对应的守护手段: - **不变量断言**:到达底层的字节数 == 被接收的 `Write` 数 × 消息长度(`TestStreamBufferWriterCloseNeverDropsRacingWrite`,200 轮竞争); - **反向验证**:临时去掉封口(注释掉 `w.closed = true`),该守护测试在 attempt 0 立即失败,证明断言真的生效而不是恒真。 ### 4.4 已排查候选与备选方案对比 在把落点收敛到流式客户端每 `Send` 两次同步交接之前,先排查了两个相邻候选:unary 的 prefix 与 body 分离写入已被上一期的 `unaryFastPathCall` 吸收成整块 body,没有剩余价值;`net.Buffers` / writev 合并 prefix 与 body 这条路的排除理由有两点,**都不是"Triple 编码没有 prefix"**:其一,**插入点的下游是 `duplexHTTPCall`,它把字节直写 `io.PipeWriter`**,而 `io.PipeWriter` 既不是 `*net.TCPConn` 也不实现 `io.ReaderFrom`,于是 writev 退化为两次顺序 `Write`,**交接次数一次未省**,还要额外付出一次 O(payload) 的拼接拷贝;其二,**真正的 socket 写发生在 x/net/http2 内部**,http2 本身就用帧缓冲把 prefix 与 body 拼进同一帧一并写出。至于 Triple 编码没有 prefix,那属于**范围说明**:这条路对它连"要合并的两段"都不存在,但它不是排除 gRPC 编码下该候选的理由。两者均不作为方案,落点� ��到流式客户端。在此基础上,对缓冲载体与触发策略做选型: | 候选 | 说明 | 结论 | | ------------------------------------ | -------------------------------------------------------- | ----------------------- | | 直接用内部 `bytes.Buffer`(构造时预分配 cap 32 KiB) | 单一流生命周期内反复复用,无归还协议 | **采用**(最小、无泄漏风险) | | 复用 `bufferPool.Get()` 但不归还 | 等同普通分配,Pool 无意义 | 放弃 | | 下沉 unary fast path 攒到 CloseWrite 才发 | 流式若攒到关闭才发,对端在流关闭前收不到数据,**破坏流式语义** | 放弃;改为按水位渐进 flush | | `net.Buffers`/writev 合并 prefix+body | 插入点下游是 `io.PipeWriter`(不是 TCPConn),writev 退化为两次顺序 Write,交接次数一次未省;且 http2 已在帧内完成 prefix 与 body 的合并 | 放弃 | ## 五、方案验证与实施 ### 5.1 改动点清单 | 文件 | 改动 | | ----------------------------------------- | ---------------------------------------------------------------------- | | `buffered_writer.go` | 新增 `streamBufferWriter`(水位 flush + `Close()` 原子封口) | | `option.go` | 新增 `WithWriteBuffering()` ClientOption(默认关) | | `protocol.go` / `client.go` | 透传 `WriteBuffering` | | `protocol_grpc.go` / `protocol_triple.go` | 流式分支装配批缓冲 + `CloseRequest` 封口刷尽且必定 `CloseWrite` | | 测试 / 基准 | `buffered_writer_test.go`、`buffered_write_bench_test.go` | | 压测工具 | `tools/benchmark/stream_ab`(流式压测客户端)、`tools/benchmark/server/dubbo-go` | ### 5.2 现有测试未覆盖的边界场景 聚合层插在协议层与 `io.Pipe` 之间,它接管的就不只是把字节送出去这一件事,还有原路径上由 `io.Pipe` 隐式兜住的一批边界。本节按边界类型归成五类,只列**现有测试覆盖不到**的场景,每一条回答三个问题:**边界条件是什么、覆盖不到会发生什么、正确行为应该是什么**。 > 体例说明:已经写出测试的场景不在此列,包括聚合语义、水位触发、超大消息直写、封口后返回 `io.ErrClosedPipe`、`Close` 幂等、空缓冲 > flush 不下沉、`Send` 与 `CloseRequest` 竞态不丢包、`Flush` 与 `Write` 并发、两条 wire 的装配与 > unary fast path 不缓冲、端到端尾包不丢,这些由 `buffered_writer_test.go` 的 10 个测试守住。本节共 13 > 条缺口,逐条的测试规划见 5.3。 #### 5.2.1 错误路径(5 条) 流是一次性资源,写失败之后这条流就废了,所以错误路径的核心要求只有两条:失败必须**可见**(不能吞成 nil),失败之后必须**可复现**(不能时好时坏)。现有测试只覆盖了「`Write` 撞上粘滞错误」这一种组合,其余出口都还是缺口。 - **首错冻结之后的另外两个出口**([buffered\_writer.go L101-L102](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L101-L102)、[L121-L122](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L121-L122)):`TestStreamBufferWriterWriteAfterError` 只断言了后续 `Write` 返回粘滞错误,没有断言 `Flush` 与 `Close` 也返回同一个错误。覆盖不到的话,「首错即冻结」只算半句话,`Close` 仍可能把已经失败的流当成功收尾。正确行为是 `err != nil` 时 `Write` / `Flush` / `Close` 三个出口一致返回同一个错误。*待补:`TestStreamBufferWriterStickyErrorOnEveryExit`* - **水位 flush 失败,但本次字节已入缓冲**([L88-L92](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L88-L92)):返回 `(len(p), nil)` 会让调用者以为这条消息发出去了、而且流还健康。正确行为是返回 `(len(p), w.err)`,字节数照实报(已被接受),错误照实报(后续不可用)。*待补:`TestStreamBufferWriterWriteReportsFlushFailure`* - **超标消息直写失败**([L81-L87](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L81-L87)):吞掉错误返回成功,调用者会以为这条大消息已经出去了。正确行为是透传底层的 `(n, err)`,不做任何包装或改写。*待补:`TestStreamBufferWriterLargeMessageErrorPropagates`* - **`Close` 时最后一次 flush 失败**([L117-L130](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L117-L130)):把失败吞成 nil,调用者会认为尾包已送达。正确行为是返回该错误,同时仍置 `closed`,不重复尝试。*待补:`TestStreamBufferWriterCloseFlushFailureReturnsError`* - **`CloseRequest` 中 flush 失败**([protocol\_grpc.go L381-L397](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/protocol_grpc.go#L381-L397)):这里有两种坏法,只返回错误而跳过 `CloseWrite`,会让对端读完请求后一直等写端关闭直到超时;两个错误都非 nil 时返回后一个,则掩盖了先发生的那次失败。正确行为是无论 flush 成败都执行 `CloseWrite`,两个错误取第一个非 nil(flush 优先)。*待补:`TestCloseRequestClosesWriteSideAfterFlushFailure`* #### 5.2.2 并发与竞态(2 条) `StreamingClientConn` 的接口契约要求 `Send` / `RequestHeader` / `CloseRequest` 三者**可互相并发**([triple.go L125-L126](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/triple.go#L125-L126)),这是 4.3 已经展开过的核心约束,此处只列还没有覆盖的两组组合。 - **`Close` 与 `Flush` 并发**([L97-L130](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L97-L130)):同一批数据会被发两次(`Flush` 在 `Close` 的尾刷之后又刷一次,或反之)。现有 `TestStreamBufferWriterConcurrent` 只并发跑了 `Write` 与 `Flush`,没有 `Close` 参与的组合。正确行为是 `closed` 判定在锁内,谁先拿到锁谁做,后到的直接返回。*待补:`TestStreamBufferWriterConcurrentCloseAndFlush`* - **并发写下的缓冲上界**([L88-L92](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L88-L92)):多个 `Send` 并发时,缓冲是否仍以 `limit` 为界、不会无界增长。现有并发测试只断言了字节总量守恒,没有断言途中任意时刻 `buf.Len() < limit`。*待补:`TestStreamBufferWriterConcurrentBufferBound`* #### 5.2.3 顺序与语义等价(2 条) 聚合层最大的风险不是丢数据,而是**改动次序**或者**少送一段字节**:这两类问题在下游看来都像偶发的奇怪故障,最难查。 - **缓冲里已有小消息,紧接着来一条 ≥ 32 KiB 的大消息**([L81-L87](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L81-L87)):大消息直写会**插队**,小消息还留在缓冲里,大消息先落了 wire,同一条流上的消息次序被打乱,接收端解出的序列与发送端不一致。现有 `TestStreamBufferWriterLargeMessageWritesDirectly` 只测了直写次数与字节数,没有「先小后大」的次序断言。正确行为是直写之前必须先把缓冲刷出去,保序责任留在聚合层内部,不要求调用者按大小排序。*待补:`TestStreamBufferWriterLargeMessageFlushesPendingFirst`* - **下游短写未被处理**([L136](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L136-L136)):`flushLocked` 只判 `err`,不看 `n < len(p)`;若下游返回「写了部分字节且 err == nil」,紧接着的 `buf.Reset()` 会把没写出去的部分丢掉,调用者完全无感。正确行为是核对 `n == buf.Len()`,短写按错误处理并冻结。*待补:`TestStreamBufferWriterShortWriteIsFailure`* #### 5.2.4 幂等与重复调用(2 条) `CloseRequest` 可能被上层在不同路径上重复触发(正常收尾、错误清理、`defer`),所以第二次调用必须明确是无事发生而不是又来一次。现有测试覆盖了「`Close` 连续两次」这一组,异常状态下的重复调用还没有覆盖。 - **`Close` 失败之后再次 `Close`**([L121-L123](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L121-L123)):返回 nil 会让调用者误以为这次关闭成功了。正确行为是在 `closed` 判定之前先判 `err`,返回粘滞错误。*待补:`TestStreamBufferWriterCloseAfterCloseFailure`* - **`Close` 之后再 `Flush`**([L104-L106](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/buffered_writer.go#L104-L106)):报 `io.ErrClosedPipe` 会把一次正常收尾判成错误。正确行为是在关闭时 flush 成功的前提下返回 nil,数据已刷尽,无事可做,也不该报错;若 flush 已经失败,则返回粘滞错误而不是 nil。*待补:`TestStreamBufferWriterFlushAfterClose`* #### 5.2.5 生命周期与装配(2 条) 这一类的边界不在聚合层内部,而在它有没有被正确接上去、以及接上去之后默认路径有没有漂移。现有装配测试覆盖了 streaming 与 unary fast path 两格,剩下两格是缺口。 - **unary 但未走 fast path**(unary fast path 的编排条件 [protocol\_grpc.go L286](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/protocol_grpc.go#L286-L307) 不成立):也会被套上缓冲。unary 只有一条消息,聚合最多省 1 次交接(Triple 编码只有一次 body 写;gRPC wire 才有 prefix 与 body 两次写),代价是一次额外拷贝加上首包推迟到 `CloseRequest`。现有测试覆盖的是 unary + fast path(不缓冲)与 streaming(缓冲)两格,unary + 非 fast path 这一格没有断言。*待补:`TestWriteBufferingCoversUnaryNonFastPath`* - **空流:一次 `Send` 都没发生就 `CloseRequest`**:请求根本没发出去,对端永远收不到请求,调用挂死。正确行为是 `duplexHTTPCall.CloseWrite` 内先 `ensureRequestMade()` 兜底发请求([duplex\_http\_call.go L130-L134](file:///home/lizining/projects/dubbo-go/protocol/triple/triple_protocol/duplex_http_call.go#L130-L134)),全仓测试对该函数 0 处引用。*待补:`TestWriteBufferingEmptyStreamStillSendsRequest`* ### 5.3 设计约束与待补测试 从上面可得以下设计约束: 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`,服务端仍能收到这次请求 | ### 5.4 端到端装配验证 单元与装配之外,还需要验证聚合层接进真实调用链之后,服务端能收到全部字节: | 测试 | 覆盖 | | ----------------------------- | ------------------------------------------------------------------------------------------ | | 单元(`buffered_writer_test.go`) | 聚合、阈值触发、超大消息直写、空 flush 空操作、错误粘滞、并发安全、`Write` after `Close` → `io.ErrClosedPipe`、`Close` 幂等 | | 并发契约守护 | `TestStreamBufferWriterCloseNeverDropsRacingWrite`(不变量 + 反向验证) | | 端到端 | `TestWriteBufferingStreamingFlushOnClose`:`CloseRequest` 刷尽、不丢尾包 | | benchmark | `buffered_write_bench_test.go`(机制层证据,ns/op / allocs) | > 说明:新增代码只在显式开启 `WithWriteBuffering()` 之后生效;默认路径下 `writeBuffer=nil`,装配给协议层的 > writer 仍是 `duplexHTTPCall`,不经过任何批缓冲逻辑。 ## 六、预估收益与方案边界 ### 6.1 生产中影响因素 1. **场景特异性(重要)**:收益集中在**小报文 × 高吞吐流式**,单流连发大量小消息、消息大小远小于缓冲水位。大报文(≥ 32 KiB 单条)直接直写底层,无收益也无损失; 2. **首包延迟**:缓冲会攒着不发,牺牲首包 latency 换吞吐,这是 opt-in 的原因; 3. **多 goroutine 并发写同一流**:聚合之后每次写都要先取锁,`w.mu` 会成为新的争用点,同一个流被多个 goroutine 并发写时需要自行权衡; 4. **端到端摊薄**:协议内收益端到端叠加框架栈/网络/服务端耗时后收敛; 5. **默认行为零变化**:未开启时 `writeBuffer=nil`,与现状逐字节等价;不适用于 unary / 大报文 / 低首包延迟场景。 ### 6.2 方案边界声明 > 这是一个**针对小报文 × 高吞吐流式场景的 opt-in 吞吐优化**,用首包延迟、一次拷贝与一把锁,换掉每条消息两次同步交接与两次 > syscall,把它们摊薄到水位级别;**默认关闭**,只对主动开启该场景的用户生效。它不改变默认行为,也不适用于所有场景。 - **适用**:流式客户端(gRPC 编码默认路径 + Triple 编码)、单流连发大量小消息、高吞吐优先; - **不适用 / 不生效**:unary、unary fast path、服务端路径、大报文(直写)、低首包延迟敏感场景; - **默认**:关闭,未开启用户零感知。 ## 七、测试构造情况说明 ### 7.1 测试构造本身是收益的变量 **收益确实与测试构造相关**:收益 = 聚合摊薄的交接次数,与**消息条数 N** 直接相关;与**单消息大小**通过 32 KiB 水位间接相关(消息越大,凑满水位的条数越少,摊薄倍数越小;单条 ≥ 32 KiB 时退化为直写)。因此测试矩阵做的是全量取值扫描,而不是只挑对优化最有利的参数组合。 ### 7.2 测试配置参照系 | 来源 | 数值 | 含义 | | --------------------------------- | ----------------------------- | ------------ | | OpenTelemetry BatchSpanProcessor | `max_export_batch_size = 512` | 每批攒 512 条再导出 | | 遥测 span 真实大小 | 200\~500 B(估算常取 500 B\~1 KB) | 单条消息量级 | | gRPC-Go / Kitex `WriteBufferSize` | 默认 32 KB | 水位一致 | | Kafka producer `batch.size` | 默认 16384(16 KB) | 跨领域同类水位 | ## 八、参考链接 - gRPC-Go `WithWriteBufferSize`(默认 32KB、Zero 关闭):<https://pkg.go.dev/google.golang.org/grpc> - gRPC C++ Performance Notes:Streaming write buffering(`set_buffer_hint` opt-in、`GRPC_ARG_HTTP2_WRITE_BUFFER_SIZE`):<http://grpc.github.io/grpc/cpp/md_doc_cpp_perf_notes.html> - gRPC-Go 官方博客 gRPC-Go performance Improvements("flush syscall for every data frame"):<https://grpc.io/blog/grpc-go-perf-improvements/> - Kitex `WithGRPCWriteBufferSize`(默认 32KB):<https://www.cloudwego.io/docs/kitex/tutorials/options/client_options/> - tRPC-Go tnet websocket 的 Combined Writes Optimization(小消息合并单 syscall):<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/> GitHub link: https://github.com/apache/dubbo-go/discussions/3673#discussioncomment-18483329 ---- This is an automatically sent email for [email protected]. To unsubscribe, please send an email to: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
