streaming_reader/
streaming_reader.rs1use std::io::{Cursor, Read};
14
15use nextjson::{from_reader, NsonDeserialize, NsonSerialize, StreamDecoder};
16
17#[derive(Debug, PartialEq, NsonSerialize, NsonDeserialize)]
19struct Message {
20 id: u64,
21 body: String,
22}
23
24struct ChunkedReader<R> {
26 inner: R,
27 chunk: usize,
28}
29
30impl<R: Read> Read for ChunkedReader<R> {
31 fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
32 let want = buf.len().min(self.chunk);
33 let mut slice = buf;
34 slice = &mut slice[..want];
35 self.inner.read(slice)
36 }
37}
38
39fn main() -> nextjson::Result<()> {
40 let messages = vec![
41 Message {
42 id: 1,
43 body: "first message".into(),
44 },
45 Message {
46 id: 2,
47 body: "second message".into(),
48 },
49 Message {
50 id: 3,
51 body: "third message".into(),
52 },
53 ];
54
55 let payload = nextjson::to_vec(&messages)?;
57 println!("序列化载荷: {} 字节", payload.len());
58
59 let slow = ChunkedReader {
60 inner: Cursor::new(payload.clone()),
61 chunk: 3,
62 };
63 let decoded: Vec<Message> = from_reader(slow)?;
64 assert_eq!(decoded, messages);
65 println!(
66 "from_reader + 3 字节/次: 解码 {} 条消息,与源数据一致",
67 decoded.len()
68 );
69
70 let single = Message {
72 id: 42,
73 body: "streamed".into(),
74 };
75 let one_payload = nextjson::to_vec(&single)?;
76 let slow_one = ChunkedReader {
77 inner: Cursor::new(one_payload),
78 chunk: 2,
79 };
80
81 let mut decoder = StreamDecoder::new(slow_one);
82 let value: Message = Message::nextdecode(&mut decoder)?;
83 decoder.end()?;
84 assert_eq!(value, single);
85 println!("StreamDecoder + 2 字节/次: 解码 msg id = {}", value.id);
86
87 let mut full = Cursor::new(nextjson::to_vec(&messages)?);
89 let mut d = StreamDecoder::new(&mut full);
90 let streamed: Vec<Message> = Vec::<Message>::nextdecode(&mut d)?;
91 d.end()?;
92 assert_eq!(streamed, messages);
93 println!(
94 "流式解码整个数组: {} 条消息,首条 id = {}",
95 streamed.len(),
96 streamed[0].id
97 );
98
99 Ok(())
100}