实现文件流阅读器批量处理并推送至RX框架
解决方案:用Rx的
Buffer实现基于数量+时间的批量处理 Great question—you’re already halfway there by leaning into Rx’s operators, and Buffer is exactly the right tool for this batching job! Let’s refine your implementation to get the exact behavior you want: batch when hitting 10 items OR after 5-10 seconds, whichever comes first.
核心思路
We’ll split responsibilities cleanly:
- Keep the file reader focused only on reading lines and pushing single messages to an Rx Subject
- Let Rx handle all the batching logic via its built-in
Bufferoperator (no manual counters or timers needed in your file watcher code)
修正后的完整实现
1. 优化文件阅读器(LogFileWatcher)
Simplify this class to only push single messages—leave batch handling to Rx:
public void StartFileWatcher(Action<LogTailMessage> callbackAction, CancellationToken cancellationToken) { var wh = new AutoResetEvent(false); var fsw = new FileSystemWatcher(_path) { Filter = _file, EnableRaisingEvents = true }; fsw.Changed += (s, e) => wh.Set(); var fileName = Path.Combine(_path, _file); var startLine = GetFileStartLine(fileName); var lineNumber = 1; using var fs = new FileStream(fileName, FileMode.Open, FileAccess.Read, FileShare.ReadWrite); using var sr = new StreamReader(fs) { while (!cancellationToken.IsCancellationRequested && !_isCancelled) { var s = sr.ReadLine(); if (s != null) { if (lineNumber >= startLine) { // Push only single messages—batching happens in Rx callbackAction(new LogTailMessage(lineNumber, s)); } lineNumber++; } else { // Wait for file changes, with a 1s timeout to avoid infinite blocking wh.WaitOne(1000); } } } }
2. 配置Rx流实现批量逻辑
Adjust the Buffer parameters to match your requirements (10 items OR 10 seconds, whichever triggers first):
var watcherSubject = new Subject<LogTailMessage>(); // Use Subject if you don't need historical messages var watcher = new LogFileWatcher(path, filename); // Start the file reader in a background task _ = Task.Run(() => watcher.StartFileWatcher(watcherSubject.OnNext, _cts.Token), _cts.Token); // Configure the Rx stream for batching Stream = watcherSubject // Core logic: emit batch when 10 items are collected OR 10 seconds pass .Buffer(TimeSpan.FromSeconds(10), 10) // Filter out empty batches (e.g., no new lines in the 10-second window) .Where(batch => batch != null && batch.Any()) // Optional: Replay last batch for new subscribers, with ref-counting to clean up when no subscribers exist .Replay(1) .RefCount();
关键细节解释
Buffer(TimeSpan, int)overload: This operator maintains an internal buffer and emits it when either:- The buffer reaches 10 items, or
- 10 seconds have passed since the last batch emission
- Subject vs ReplaySubject: Use
Subjectif new subscribers don’t need to see historical messages (lighter weight). Stick withReplaySubjectonly if you need to retain past batches for late joiners, and limit the replay buffer size to avoid memory leaks. - FileSystemWatcher Debouncing: The
Changedevent can fire multiple times for a single file write. To avoid unnecessary wake-ups, you could add a simple debounce (e.g., usingTask.Delayand a cancellation token) before callingwh.Set().
内容的提问来源于stack exchange,提问作者NZJames
相关产品推荐
相关产品推荐

