如何以类似serde-jsonlines的方式并行读取bincode序列化结构体?
问题与解答
问题
我目前使用serde-jsonlines将大量相同结构体序列化到文件,再通过rayon的par_bridge并行读取,代码如下:
let mut reader = JsonLinesReader::new(input_file); let results: Vec<ResultStruct> = db_json_reader .read_all::<MyStruct>() .par_bridge() .into_par_iter() .map(|my_struct| { // 处理结构体并返回结果 result }) .collect();
该方案可行是因为JsonLinesReader会返回按行分割的迭代器。我希望改用bincode编码结构体以减小文件体积,已实现如下可行的序列化与反序列化示例:
use bincode; use serde::{Deserialize, Serialize}; use std::fs::File; use std::io::{BufWriter, Write}; #[derive(Debug, Deserialize, Serialize)] struct MyStruct { name: String, value: Vec<u64>, } pub fn playground() { let s1 = MyStruct { name: "Hello".to_string(), value: vec![1, 2, 3], }; let s2 = MyStruct { name: "World!".to_string(), value: vec![3, 4, 5, 6], }; let out_file = File::create("test.bin").expect("Unable to create file"); let mut writer = BufWriter::new(out_file); let s1_encoded: Vec<u8> = bincode::serialize(&s1).unwrap(); writer.write_all(&s1_encoded).expect("Unable to write data"); let s2_encoded: Vec<u8> = bincode::serialize(&s2).unwrap(); writer.write_all(&s2_encoded).expect("Unable to write data"); drop(writer); let mut in_file = File::open("test.bin").expect("Unable to open file"); let s1_decoded: MyStruct = bincode::deserialize_from(&mut in_file).expect("Unable to read data"); let s2_decoded: MyStruct = bincode::deserialize_from(&mut in_file).expect("Unable to read data"); println!("s1_decoded: {:?}", s1_decoded); println!("s2_decoded: {:?}", s2_decoded); }
请问是否可以采用类似serde-jsonlines的方式并行读取这些结构体?我猜测可能无法实现,因为每个结构体并非以换行符分隔,无法合理拆分输入流以支持多线程处理。
回答
你的猜测是正确的,直接使用当前的bincode序列化方式无法像serde-jsonlines那样实现并行读取,核心原因是序列化流中没有明确的结构体分隔标记,bincode的输出是连续字节流,无法在任意位置拆分并独立解析单个结构体。
不过可以通过修改序列化格式来实现并行读取,以下是两种可行方案:
1. 手动添加长度前缀
在序列化每个结构体前,先写入该结构体序列化后的字节长度(用固定大小的数值类型,比如u64),这样读取时可以先获取长度,再精准截取对应字节段进行解析。
修改后的序列化代码:
pub fn serialize_with_length_prefix() { let s1 = MyStruct { name: "Hello".to_string(), value: vec![1, 2, 3], }; let s2 = MyStruct { name: "World!".to_string(), value: vec![3, 4, 5, 6], }; let out_file = File::create("test_with_prefix.bin").expect("无法创建文件"); let mut writer = BufWriter::new(out_file); // 序列化并写入s1:先写长度,再写数据 let s1_encoded = bincode::serialize(&s1).unwrap(); let len = s1_encoded.len() as u64; bincode::serialize_into(&mut writer, &len).unwrap(); writer.write_all(&s1_encoded).unwrap(); // 序列化并写入s2:先写长度,再写数据 let s2_encoded = bincode::serialize(&s2).unwrap(); let len = s2_encoded.len() as u64; bincode::serialize_into(&mut writer, &len).unwrap(); writer.write_all(&s2_encoded).unwrap(); }
并行读取实现思路:
- 预扫描文件偏移:遍历文件,依次读取每个
u64长度值,计算出每个结构体数据的起始偏移(当前位置 + 8字节,因为u64占8字节)和数据长度,将这些偏移信息收集为一个列表。 - 并行解析:使用
rayon对偏移列表做并行迭代,每个线程根据偏移信息打开文件(或通过内存映射mmap高效访问),截取对应字节段后解析为MyStruct并处理。
2. 使用bincode自带的长度帧配置
bincode的Options提供了with_length_prefix方法,可以自动为每个序列化的结构体添加长度前缀,简化序列化流程:
use bincode::Options; pub fn serialize_with_bincode_frame() { let s1 = MyStruct { name: "Hello".to_string(), value: vec![1, 2, 3], }; let s2 = MyStruct { name: "World!".to_string(), value: vec![3, 4, 5, 6], }; let out_file = File::create("test_with_frame.bin").expect("无法创建文件"); let mut writer = BufWriter::new(out_file); let options = bincode::DefaultOptions::new().with_length_prefix(); options.serialize_into(&mut writer, &s1).unwrap(); options.serialize_into(&mut writer, &s2).unwrap(); }
这种方式的并行读取思路和第一种方案一致,先预扫描所有帧的位置,再并行处理每个帧的数据。
优化建议
- 处理大文件时,使用
memmap2crate将文件映射到内存,避免每个线程重复打开文件或频繁seek,提升读取效率。 - 若结构体数量极大,预扫描偏移的过程会有一定开销,但后续并行处理的性能收益通常能抵消这部分成本。
内容的提问来源于stack exchange,提问作者DHP
相关产品推荐
相关产品推荐

