Client Streaming
หลาย request เข้า, 1 response ออก
หัวข้อที่มีชื่อว่า “หลาย request เข้า, 1 response ออก”Client streaming กลับด้านจาก server streaming: client ส่ง ลำดับ ของ message บน stream ที่เปิดค้างไว้ แล้ว server ตอบกลับด้วย response เดียว — แต่หลังจาก client บอกว่าส่งเสร็จแล้วเท่านั้น
รูปแบบนี้เหมาะกับการป้อน data เข้า service: upload ไฟล์เป็น chunk, รับ metric เป็น batch หรือส่งหลาย record ที่ server รวบเป็นสรุปเดียว
sequenceDiagram participant C as Client participant S as Server C->>S: Sample 1 C->>S: Sample 2 C->>S: Sample 3 C->>S: (half-close: ส่งเสร็จแล้ว) S-->>C: UploadSummary(count: 3, accepted: 3) Note over C,S: server ตอบครั้งเดียว หลัง client ส่งเสร็จ
contract
หัวข้อที่มีชื่อว่า “contract”ตรงนี้ keyword stream อยู่ที่ request ส่วน response ยังเป็นตัวเดียว:
syntax = "proto3";package metrics.v1;
message Sample { string name = 1; double value = 2;}
message UploadSummary { int32 count = 1; int32 accepted = 2;}
service Ingest { // A stream of Samples, one UploadSummary response. rpc UploadSamples(stream Sample) returns (UploadSummary);}Client: ส่งใน loop แล้วปิด
หัวข้อที่มีชื่อว่า “Client: ส่งใน loop แล้วปิด”client ส่งแต่ละ message แล้วทำ half-close — บอก server ว่า “ส่งเสร็จแล้ว” — แล้วรอ response เดียว
stream, _ := client.UploadSamples(ctx)for _, s := range samples { if err := stream.Send(s); err != nil { log.Fatal(err) }}summary, err := stream.CloseAndRecv() // half-close, then get the replyif err != nil { log.Fatal(err)}fmt.Println(summary.Accepted)def sample_iter(): for s in samples: yield metrics_pb2.Sample(name=s.name, value=s.value)
summary = client.UploadSamples(sample_iter()) # returns after the iterator is exhaustedprint(summary.accepted)const call = client.uploadSamples((err, summary: UploadSummary) => { if (err) throw err; console.log(summary.accepted);});for (const s of samples) call.write(s);call.end(); // half-close: done sending, now wait for the responseServer: ยุบ stream ให้เหลือคำตอบเดียว
หัวข้อที่มีชื่อว่า “Server: ยุบ stream ให้เหลือคำตอบเดียว”server อ่านทุก message, สะสม state, แล้ว return response เดียวเมื่อ client half-close
func (s *server) UploadSamples(stream metricsv1.Ingest_UploadSamplesServer) error { var count, accepted int32 for { sample, err := stream.Recv() if err == io.EOF { // client is done — send the single response return stream.SendAndClose(&metricsv1.UploadSummary{Count: count, Accepted: accepted}) } if err != nil { return err } count++ if sample.Value >= 0 { accepted++ } }}def UploadSamples(self, request_iterator, context): count = accepted = 0 for sample in request_iterator: count += 1 if sample.value >= 0: accepted += 1 return metrics_pb2.UploadSummary(count=count, accepted=accepted)function uploadSamples(call: ServerReadableStream<Sample, UploadSummary>, cb) { let count = 0, accepted = 0; call.on('data', (s: Sample) => { count++; if (s.value >= 0) accepted++; }); call.on('end', () => cb(null, { count, accepted }));}หมายเหตุด้านการออกแบบ
หัวข้อที่มีชื่อว่า “หมายเหตุด้านการออกแบบ”- response มาหลัง half-close เท่านั้น client อย่าคาดหวังคำตอบกลาง stream response เดียวคือสัญญาณว่า upload ทั้งก้อนถูกประมวลผลแล้ว
- จำกัดยอดรวม client stream ได้ไม่รู้จบ server ควรจำกัดว่าจะรับแค่ไหน (จำนวน message หรือขนาด) เพื่อกัน memory หรืองานที่ไม่จำกัด
- ลำดับถูกรักษาไว้ server จึงพึ่งได้ว่าได้รับ sample ตามลำดับที่ส่ง
- ทั้งสองฝั่งจบก่อนด้วย error ได้ ถ้า server ปฏิเสธ batch กลางทาง (เช่น quota เกิน) server ก็ return error status แทน summary ได้ —
CloseAndRecvของ client จะเจอ error นั้น