使用System.Text.Json实现非均匀JSON数据的流式解析方案
使用System.Text.Json流式解析非均匀SPARQL JSON结果的最简方法
需要流式解析SPARQL JSON格式的结果数据,这类数据的顶层结构是对象而非数组,因此无法直接使用JsonSerializer.DeserializeAsyncEnumerable<T>。数据来自HTTP流,准长时间运行的查询需要将结果增量集成到客户端UI以提升体验。
SPARQL JSON示例:
{ "head": { "vars": [ "s" , "p" , "o" ] } , "results": { "bindings": [ { "s": { "type": "uri" , "value": "https://example.com/ontology/al/bnode/1" } , "p": { "type": "uri" , "value": "http://www.w3.org/1999/02/22-rdf-syntax-ns#type" } , "o": { "type": "uri" , "value": "https://example.com/ontology/al/Module" } } , { "s": { "type": "bnode" , "value": "b0" } , "p": { "type": "uri" , "value": "http://www.w3.org/1999/02/22-rdf-syntax-ns#rest" } , "o": { "type": "uri" , "value": "http://www.w3.org/1999/02/22-rdf-syntax-ns#nil" } } , { "s": { "type": "triple" , "value": { "subject": { "type": "uri" , "value": "https://example.com/ontology/al/bnode/1798" } , "predicate": { "type": "uri" , "value": "https://example.com/ontology/al/type" } , "object": { "type": "uri" , "value": "https://example.com/ontology/al/bnode/6" } } } , "p": { "type": "uri" , "value": "https://example.com/ontology/al/dimensions" } , "o": { "type": "bnode" , "value": "b0" } } , { "s": { "type": "uri" , "value": "https://example.com/ontology/al/bnode/1" } , "p": { "type": "uri" , "value": "https://example.com/ontology/al/publisher" } , "o": { "type": "literal" , "value": "Example" } } ] } }
其中head/vars定义了查询的变量集合,results/bindings数组中的每个对象对应一组变量绑定值。已定义的解析类型如下(F#):
[<Struct>] type RdfUri = { Value: string } member this.Uri = Uri(this.Value) [<Struct>] type RdfBlankNode = { Value: string } [<Struct>] type RdfLiteral = { Value: string DataType: string voption } type RdfTriple = { Subject: RdfValue Predicate: RdfValue Object: RdfValue } [<RequireQualifiedAccess>] type RdfValue = | Uri of RdfUri | BlankNode of RdfBlankNode | Literal of RdfLiteral /// RDF-star | Triple of RdfTriple // SPARQL types type SparqlVariableBinding = { Name: string Value: RdfValue } type SparqlVariableBindings = { Bindings: block<SparqlVariableBinding> } member this.TryFind(name: string) : RdfValue option = this.Bindings |> Block.tryFind (fun var -> String.Equals(var.Name, name, StringComparison.Ordinal)) |> Option.map (fun var -> var.Value) member this.Item(name: string) : RdfValue = match this.TryFind name with | None -> failwith $"Could not find variable %s{name}" | Some x -> x type SparqlHead = { VariableNames: Set<string> } [<RequireQualifiedAccess>] type SparqlResult = | Head of SparqlHead | Variables of SparqlVariableBindings [<RequireQualifiedAccess>] type AsyncSparqlResultSet = AsyncSeq<SparqlResult>
最简流式解析方案
核心思路是使用Utf8JsonReader手动遍历JSON令牌,定位到目标节点后增量解析,避免一次性加载整个JSON到内存。
1. 初始化异步流式读取器
open System.Text.Json let createJsonReader (stream: Stream) = let options = JsonReaderOptions(AllowTrailingCommas = true, CommentHandling = JsonCommentHandling.Skip) Utf8JsonReader(stream, options)
2. 解析Head节点
遍历令牌找到"head"对象,解析其中的"vars"数组得到变量名集合:
let rec parseHead (reader: byref<Utf8JsonReader>) = seq { while reader.Read() do match reader.TokenType with | JsonTokenType.PropertyName when reader.ValueTextEquals("vars") -> reader.Read() // 进入数组 let vars = ResizeArray<string>() while reader.Read() && reader.TokenType <> JsonTokenType.EndArray do vars.Add(reader.GetString()) yield SparqlResult.Head { VariableNames = Set vars } return () | JsonTokenType.EndObject -> return () | _ -> () }
3. 解析RDF值与绑定对象
先实现RdfValue的递归解析,支持RDF-star的triple类型:
let rec parseRdfValue (reader: byref<Utf8JsonReader>) = let rec parseTriple (reader: byref<Utf8JsonReader>) = let mutable subject = Unchecked.defaultof<RdfValue> let mutable predicate = Unchecked.defaultof<RdfValue> let mutable obj = Unchecked.defaultof<RdfValue> while reader.Read() do match reader.TokenType with | JsonTokenType.PropertyName when reader.ValueTextEquals("subject") -> reader.Read() subject <- parseRdfValue &reader | JsonTokenType.PropertyName when reader.ValueTextEquals("predicate") -> reader.Read() predicate <- parseRdfValue &reader | JsonTokenType.PropertyName when reader.ValueTextEquals("object") -> reader.Read() obj <- parseRdfValue &reader | JsonTokenType.EndObject -> return { Subject = subject; Predicate = predicate; Object = obj } | _ -> () failwith "Invalid triple structure" let mutable rdfType = "" let mutable rdfValue = Unchecked.defaultof<JsonElement> while reader.Read() do match reader.TokenType with | JsonTokenType.PropertyName when reader.ValueTextEquals("type") -> reader.Read() rdfType <- reader.GetString() | JsonTokenType.PropertyName when reader.ValueTextEquals("value") -> reader.Read() rdfValue <- reader.Clone() | JsonTokenType.EndObject -> match rdfType with | "uri" -> RdfValue.Uri { Value = rdfValue.GetString() } | "bnode" -> RdfValue.BlankNode { Value = rdfValue.GetString() } | "literal" -> let dataType = if reader.Read() && reader.TokenType = JsonTokenType.PropertyName && reader.ValueTextEquals("datatype") then reader.Read() Some(reader.GetString()) |> ValueSome else ValueNone RdfValue.Literal { Value = rdfValue.GetString(); DataType = dataType } | "triple" -> let mutable tripleReader = rdfValue.GetReader() let triple = parseTriple &tripleReader RdfValue.Triple triple | _ -> failwith $"Unknown RDF type: {rdfType}" |> return | _ -> () failwith "Invalid RDF value structure"
再解析单个变量绑定对象:
let parseBinding (reader: byref<Utf8JsonReader>) = let bindings = ResizeArray<SparqlVariableBinding>() while reader.Read() do match reader.TokenType with | JsonTokenType.PropertyName -> let varName = reader.GetString() reader.Read() // 进入值对象 let rdfValue = parseRdfValue &reader bindings.Add({ Name = varName; Value = rdfValue }) | JsonTokenType.EndObject -> { Bindings = Block.ofSeq bindings } |> SparqlResult.Variables |> return | _ -> () failwith "Invalid binding structure"
4. 组合完整异步解析流程
将上述步骤整合为返回AsyncSeq<SparqlResult>的函数,实现增量输出:
open FSharp.Control let parseSparqlResultsAsync (stream: Stream) : AsyncSparqlResultSet = asyncSeq { use stream = stream let mutable reader = createJsonReader stream // 输出Head结果 yield! parseHead &reader |> AsyncSeq.ofSeq // 定位到bindings数组并增量解析 while reader.Read() do match reader.TokenType with | JsonTokenType.PropertyName when reader.ValueTextEquals("results") -> reader.Read() // 进入results对象 while reader.Read() do match reader.TokenType with | JsonTokenType.PropertyName when reader.ValueTextEquals("bindings") -> reader.Read() // 进入bindings数组 while reader.Read() && reader.TokenType <> JsonTokenType.EndArray do yield parseBinding &reader return () | _ -> () return () | _ -> () }
方案说明
- 该方案通过
Utf8JsonReader实现真正的流式解析,无需加载完整JSON到内存,适配大体积或长时间运行的SPARQL查询。 - 解析过程先输出
Head结果,随后增量输出每个Variables绑定,可直接用于客户端UI实时更新。 - 原生支持RDF-star的
triple类型递归解析,适配非均匀JSON结构。
内容的提问来源于stack exchange,提问作者Bent Rasmussen
相关产品推荐
相关产品推荐

