java中实现游标分批查询与异步入库流水线,核心是解耦读写、用主键游标避免深度分页、blockingqueue缓冲、每批100–1000条事务批量入库、checkpoint持久化支持断点续传。

Java 中实现基于游标的分批查询与批量异步入库流水线,核心是解耦“读取”和“写入”,用游标避免全表扫描与重复数据,用异步+缓冲提升吞吐,同时保证顺序性与失败可恢复。关键不在单个技术点,而在各环节的衔接控制。
游标分页查询:用 last_id 或时间戳做安全断点
替代传统 limit/offset(深度分页性能差),每次查出一批后记录最后一条主键或更新时间,作为下一次查询条件。推荐用主键(如 id)——单调递增、无重复、索引友好。
- SQL 示例:SELECT * FROM order WHERE id > ? ORDER BY id LIMIT 1000
- 首次查询用 id > 0 或直接 id >= 1;后续传入上一批最大 id
- 务必按游标字段 严格升序 + 索引覆盖,否则可能漏数或重复
- 查询结果为空时,说明数据读完,流程自然终止
生产者-消费者模型:用 BlockingQueue 缓冲批次
查询线程(生产者)持续拉取游标批次,转成 List
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
- 选 LinkedBlockingQueue,设合理容量(如 10–50),防内存溢出
- 生产者不等待消费,但需处理队列满时的策略:可阻塞(默认)、丢弃、或主动限流(如 Thread.sleep)
- 消费者建议单线程(保序)或多线程(需按业务键分片,如 userId % N,避免同一条记录被乱序写入)
异步批量入库:JDBC 批处理 + 事务粒度控制
每批次调用一次 PreparedStatement.addBatch() + executeBatch(),比单条插入快 5–20 倍。事务不宜过大(防锁表、OOM),也不宜过小(降低吞吐)。
- 推荐每批 100–1000 条,具体看单行大小与数据库承受力
- 每个批次一个独立事务:成功则提交,失败则整批重试或落日志告警,不污染其他批次
- 使用 JdbcTemplate.batchUpdate() 或 MyBatis 的 foreach + useGeneratedKeys="false" 关闭主键回填以提速
- 注意连接池配置:maxActive 要够,避免消费者抢不到连接而阻塞
容错与可观测:断点续传 + 进度跟踪
流水线长期运行必须支持中断恢复,不能每次从头开始。
- 每次成功入库一批后,把当前游标值(如 max_id)持久化到 DB 或 Redis,而非仅存在内存
- 启动时优先查这个“checkpoint”,有则从它继续;无则从初始游标开始
- 加简单监控:打印已处理总条数、当前游标、队列剩余量、最近批次耗时,可用 AtomicLong + 定时日志
- 对写入失败批次,记录错误原因和原始数据(可截取前几条),便于人工干预或自动重推
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南










