如何用serde_json构建有状态的流式JSON解析器?
有状态流式JSON解析(Serde/Serde_json)的优化方案
需求说明
- 处理超大体积JSON,无法全量加载到内存,需流式解析;同时因JSON嵌套层级深,需使用
disable_recursion_limit关闭递归限制 - 支持有状态处理,通过传入的状态数据决定输入JSON内容的保留规则与转换逻辑
示例输入JSON:
{ "documents": [ { "foo": 1 }, { "baz": true }, { "bar": null } ], "journal": { "timestamp": "2023-04-04T08:28:00" } }
需重点处理documents数组(元素体积大,仅保留子集),其他键暂不关注。
当前实现的问题
- 无法流式处理
documents数组,不知道如何将状态传递给子解析器 - 代码结构冗余冗长,不符合Rust惯例
- 尝试直接用
MapAccess调用DeserializeSeed::deserialize时,因类型不匹配报错
当前实现代码:
use serde::de::DeserializeSeed; use serde_json::Value; /// A simplified state passed to and returned from the serialization. #[derive(Debug, Default)] struct Stats { records_skipped: usize, } /// Models the input data; `Documents` is just a vector of JSON values, /// but it is its own type to allow custom deserialization #[derive(Debug)] struct MyData { documents: Vec<Value>, journal: Value, } struct MyDataDeserializer<'a> { state: &'a mut Stats, } /// Top-level seeded deserializer only so I can plumb the state through impl<'de> DeserializeSeed<'de> for MyDataDeserializer<'_> { type Value = MyData; fn deserialize<D>(mut self, deserializer: D) -> Result<Self::Value, D::Error> where D: serde::Deserializer<'de>, { let visitor = MyDataVisitor(&mut self.state); let docs = deserializer.deserialize_map(visitor)?; Ok(docs) } } struct MyDataVisitor<'a>(&'a mut Stats); impl<'de> serde::de::Visitor<'de> for MyDataVisitor<'_> { type Value = MyData; fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result { write!(formatter, "a map") } fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error> where A: serde::de::MapAccess<'de>, { let mut documents = Vec::new(); let mut journal = Value::Null; while let Some(key) = map.next_key::<String>()? { println!("Got key = {key}"); match &key[..] { "documents" => { // Not sure how to handle the next value in a streaming manner documents = map.next_value()?; } "journal" => journal = map.next_value()?, _ => panic!("Unexpected key '{key}'"), } } Ok(MyData { documents, journal }) } } struct DocumentDeserializer<'a> { state: &'a mut Stats, } impl<'de> DeserializeSeed<'de> for DocumentDeserializer<'_> { type Value = Vec<Value>; fn deserialize<D>(mut self, deserializer: D) -> Result<Self::Value, D::Error> where D: serde::Deserializer<'de>, { let visitor = DocumentVisitor(&mut self.state); let documents = deserializer.deserialize_seq(visitor)?; Ok(documents) } } struct DocumentVisitor<'a>(&'a mut Stats); impl<'de> serde::de::Visitor<'de> for DocumentVisitor<'_> { type Value = Vec<Value>; fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result { write!(formatter, "a list") } fn visit_seq<A>(self, mut seq: A) -> Result<Self::Value, A::Error> where A: serde::de::SeqAccess<'de>, { let mut agg_map = serde_json::Map::new(); while let Some(item) = seq.next_element()? { // If `item` isn't a JSON object, we'll skip it: let Value::Object(map) = item else { continue }; // Get the first element, assuming we have some let (k, v) = match map.into_iter().next() { Some(kv) => kv, None => continue, }; // Ignore any null values; aggregate everything into a single map if v == Value::Null { self.0.records_skipped += 1; continue; } else { println!("Keeping {k}={v}"); agg_map.insert(k, v); } } let values = Value::Object(agg_map); println!("Final value is {values}"); Ok(vec![values]) } } fn main() { let fh = std::fs::File::open("input.json").unwrap(); let buf = std::io::BufReader::new(fh); let read = serde_json::de::IoRead::new(buf); let mut state = Stats::default(); let mut deserializer = serde_json::Deserializer::new(read); let mydata = MyDataDeserializer { state: &mut state } .deserialize(&mut deserializer) .unwrap(); println!("{mydata:?}"); }
优化后的解决方案
核心思路
利用Serde的DeserializeSeed结合MapAccess的next_value_seed方法,直接在父解析器中传递状态给子解析器,保持流式处理特性;通过简化Visitor与解析器结构,让代码更简洁符合Rust惯例。
完整代码实现
use serde::{de::DeserializeSeed, Deserialize}; use serde_json::{self, Value}; use std::fmt; #[derive(Debug, Default)] struct Stats { records_skipped: usize, } #[derive(Debug)] struct MyData { documents: Value, // 直接存储聚合后的Object,避免冗余Vec包装 journal: Value, } // 子解析器:处理documents数组,携带状态 struct DocumentsSeed<'a> { state: &'a mut Stats, } impl<'de> DeserializeSeed<'de> for DocumentsSeed<'_> { type Value = Value; fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> where D: serde::Deserializer<'de>, { deserializer.deserialize_seq(DocumentsVisitor(self.state)) } } struct DocumentsVisitor<'a>(&'a mut Stats); impl<'de> serde::de::Visitor<'de> for DocumentsVisitor<'_> { type Value = Value; fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { write!(formatter, "a sequence of objects") } fn visit_seq<A>(self, mut seq: A) -> Result<Self::Value, A::Error> where A: serde::de::SeqAccess<'de>, { let mut agg_map = serde_json::Map::new(); while let Some(obj) = seq.next_element::<Value>()? { let Value::Object(map) = obj else { continue }; if let Some((k, v)) = map.into_iter().next() { if v == Value::Null { self.0.records_skipped += 1; } else { agg_map.insert(k, v); } } } Ok(Value::Object(agg_map)) } } // 顶层解析器:处理整个JSON对象,携带状态 struct MyDataSeed<'a> { state: &'a mut Stats, } impl<'de> DeserializeSeed<'de> for MyDataSeed<'_> { type Value = MyData; fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> where D: serde::Deserializer<'de>, { deserializer.deserialize_map(MyDataVisitor(self.state)) } } struct MyDataVisitor<'a>(&'a mut Stats); impl<'de> serde::de::Visitor<'de> for MyDataVisitor<'_> { type Value = MyData; fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { write!(formatter, "a top-level JSON object") } fn visit_map<A>(mut self, mut map: A) -> Result<Self::Value, A::Error> where A: serde::de::MapAccess<'de>, { let mut documents = Value::Null; let mut journal = Value::Null; while let Some(key) = map.next_key::<String>()? { match key.as_str() { "documents" => { // 使用next_value_seed传递状态,实现流式解析 documents = map.next_value_seed(DocumentsSeed { state: self.0 })?; } "journal" => { journal = map.next_value()?; } _ => { // 跳过未知键,增强代码健壮性 map.next_value::<Value>()?; } } } Ok(MyData { documents, journal }) } } fn main() { let fh = std::fs::File::open("input.json").unwrap(); let buf = std::io::BufReader::new(fh); let mut state = Stats::default(); let mut deserializer = serde_json::Deserializer::new(buf) .disable_recursion_limit(); // 启用递归限制关闭 let mydata = MyDataSeed { state: &mut state } .deserialize(&mut deserializer) .unwrap(); println!("Parsed data: {mydata:?}"); println!("Stats: {state:?}"); }
关键优化点
- 使用
next_value_seed替代next_value,直接在MapAccess中传递状态给子解析器,解决类型不匹配问题,同时保持流式解析特性 - 简化
MyData结构,将聚合后的结果直接存为Value::Object,避免不必要的Vec<Value>包装 - 跳过未知键而非panic,增强代码健壮性
- 合并冗余的解析器逻辑,让代码结构更简洁符合Rust风格
- 正确启用
disable_recursion_limit,适配深层嵌套的JSON结构
内容的提问来源于stack exchange,提问作者user655321
相关产品推荐
相关产品推荐

