如何在Rust程序中实时读取每5分钟追加数据的CSV文件?
解决方案
核心思路
不用全量重读CSV,而是记录上次读取的位置(文件偏移量),或者利用CSV自带的序号/时间戳做增量校验,结合notify监听文件变更事件,只读取新增行即可。
关于防抖器的必要性
你的Python脚本是每5分钟追加一次,属于低频、单次的文件变更(每次追加是一次性写入),这种场景下不需要防抖器——不会出现短时间内频繁触发变更事件的情况。但如果Python脚本是分批多次写入(比如1分钟内多次追加),可以加个1秒左右的简单延迟防抖,避免重复处理。
具体实现步骤
1. 初始化:记录初始读取状态
程序启动时先全量读取CSV,同时记录两个关键状态:
- 文件的当前偏移量:可以通过
std::fs::metadata获取文件大小,或者用csv::Reader的byte_offset()方法 - 最后一行的序号/时间戳:作为兜底校验,防止文件被意外修改导致偏移量失效
2. 用notify监听文件变更
使用notify的RecommendedWatcher监听CSV文件的Write或Modify事件,触发时执行增量读取。
3. 增量读取的两种方式
方式一:基于文件偏移量(高效)
直接从上次记录的偏移量开始读取文件内容,解析成CSV行:
use std::fs::File; use std::io::{BufReader, Seek, SeekFrom}; use csv::ReaderBuilder; // 假设last_offset是上次记录的文件偏移量 let mut file = File::open("data.csv")?; file.seek(SeekFrom::Start(last_offset))?; let mut rdr = ReaderBuilder::new().from_reader(BufReader::new(file)); for result in rdr.records() { let record = result?; // 处理新增记录逻辑 } // 更新last_offset为当前文件大小 let new_offset = file.metadata()?.len();
方式二:基于序号/时间戳(容错性高)
如果担心文件被截断、历史行被修改,用CSV里的序号或时间戳过滤,只处理新增行:
use csv::ReaderBuilder; use serde::Deserialize; #[derive(Debug, Deserialize)] struct Record { timestamp: String, seq: u32, // 其他业务字段 } // 假设last_seq是上次记录的最大序号 let mut rdr = ReaderBuilder::new().from_path("data.csv")?; for result in rdr.deserialize() { let record: Record = result?; if record.seq > last_seq { // 处理新增记录 last_seq = record.seq; } }
这种方式会遍历全文件,但因为有序号过滤,实际只处理新增行,适合文件体积不大的场景。
4. 异常处理
- 监听文件被删除、重命名的事件,此时需要重新初始化读取状态
- 处理CSV解析错误(比如Python脚本写入不完整的行),可以延迟几秒再重试读取
简化版完整代码
use notify::{RecommendedWatcher, RecursiveMode, Watcher}; use std::fs::File; use std::io::{BufReader, Seek, SeekFrom}; use std::path::Path; use std::sync::mpsc::channel; use std::time::Duration; fn main() -> Result<(), Box<dyn std::error::Error>> { // 初始化:读取全量数据并记录初始偏移量 let mut last_offset = 0; let path = Path::new("data.csv"); if path.exists() { let file = File::open(path)?; last_offset = file.metadata()?.len(); // 此处添加全量读取CSV的逻辑 } // 创建监听通道 let (tx, rx) = channel(); let mut watcher = RecommendedWatcher::new(tx, Duration::from_secs(2))?; watcher.watch(path, RecursiveMode::NonRecursive)?; println!("Watching CSV file for changes..."); loop { match rx.recv() { Ok(event) => match event.kind { notify::EventKind::Modify(_) => { println!("File modified, reading new records..."); let mut file = File::open(path)?; // 跳转到上次读取的位置 file.seek(SeekFrom::Start(last_offset))?; let mut rdr = csv::ReaderBuilder::new().from_reader(BufReader::new(file)); for result in rdr.records() { match result { Ok(record) => println!("New record: {:?}", record), Err(e) => eprintln!("Failed to parse record: {}", e), } } // 更新偏移量为当前文件大小 last_offset = file.metadata()?.len(); } _ => {} // 忽略删除、创建等其他事件 }, Err(e) => eprintln!("Watch error: {}", e), } } }
内容的提问来源于stack exchange,提问作者snazzybeaver
相关产品推荐
相关产品推荐

