并发缓冲容器本身不自动记断点,而是作为线程安全的断点暂存中转站,需配合主动行号维护、异常捕获和最终落盘三步闭环;Java 中常用 AtomicInteger 或 ConcurrentHashMap 存储当前行号,其中 AtomicInteger 最轻量且支持原子操作;初始化从断点文件读取初始值,每成功读一行即 incrementAndGet() 更新,finally 中持久化最新行号;多线程按行分配而非字节切分,主线程预生成行偏移表并用 ConcurrentLinkedQueue 分发,各线程独立维护 localLineNum;重启时以断点文件为准初始化 AtomicInteger,跳过已处理行时不计入业务逻辑且不触发保存。

用并发缓冲容器在本地文件读取过程中平滑存储中断点,核心不是靠容器本身“自动记断点”,而是把容器作为线程安全的断点暂存中转站,配合主动行号维护、异常捕获和最终落盘三步闭环。Java 中没有专为“断点存储”设计的并发容器,但可用 AtomicInteger 或 ConcurrentHashMap<string integer></string>(如按文件名索引)来承载正在读取的当前行号,确保多线程环境下计数不丢失、不覆盖。
用 AtomicInteger 做线程安全的行号计数器
这是最轻量也最常用的做法。它天然支持原子递增、获取和设置,避免 synchronized 块或锁带来的开销。
- 初始化时从断点文件读取初始值(如 127),设为
AtomicInteger checkpoint = new AtomicInteger(initialLine) - 每次
readLine()成功后立即调用checkpoint.incrementAndGet(),保证行号与实际处理严格同步 - 不要在循环变量或 for 索引上做文章——那些不反映真实读取进度
- 注意:incrementAndGet 返回的是递增后的值,正好对应“刚处理完的这一行”的行号
在 finally 块中把最新行号写入持久化文件
并发容器只管内存中的暂存,真正“平滑”靠的是异常发生时仍能落地。无论业务解析崩了、JSON 格式错,还是 IO 被中断,都要确保最后已确认的行号被保存。
- 整个读取逻辑包在 while 循环内,每行处理用 try-catch 包裹
- 该 try 块外加一个 finally,只做一件事:
saveCheckpoint(checkpoint.get()) -
saveCheckpoint()内部推荐用Files.writeString(path, String.valueOf(lineNum)),避开流管理,简洁可靠 - 如果 saveCheckpoint 自身可能失败(如磁盘满),可加一层重试或降级日志,但不能让它阻塞主流程
多线程读同一文件时,按行分配而非按字节切分
并发缓冲容器解决不了“跨行读取”问题。若多个线程直接 seek 到随机偏移读,极易把一行切成两半。必须把“行”作为调度单元。
- 主线程先扫描一次文件,用
RandomAccessFile或Files.lines()(小文件)收集所有换行符位置,生成行偏移表 - 用
ConcurrentLinkedQueue<long></long>存这些起始偏移,每个线程 poll 一个位置,fseek后用BufferedReader读完整行 - 每线程各自维护自己的
AtomicInteger localLineNum,最终汇总或单独落盘,避免竞争 - 禁止让线程共享同一个 BufferedReader 实例——它不是线程安全的
断点恢复时跳过已处理行,不依赖容器状态
并发容器里的数值只在运行时有效,重启后清空。真正恢复依据是外部断点文件,不是内存里的 AtomicInteger。
- 启动时优先读
checkpoint.txt,若不存在则从 1 开始;若存在,就用该值初始化 AtomicInteger - 恢复读取时,不调用
skipNLines()(不可靠),而是新建 BufferedReader,循环readLine()直到达到目标行号 - 跳过过程不计入业务逻辑,也不触发 saveCheckpoint,只算初始化动作
- 断点文件保持纯文本整数格式,不引入 JSON 或版本字段,降低损坏风险

















