You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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数组(元素体积大,仅保留子集),其他键暂不关注。

当前实现的问题

  1. 无法流式处理documents数组,不知道如何将状态传递给子解析器
  2. 代码结构冗余冗长,不符合Rust惯例
  3. 尝试直接用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 03:44:54