Java数据采集系统实战:从阻塞队列到增量同步的完整链路
2026/9/14 10:57:20 网站建设 项目流程

简介:这是一个基于Java实现的数据采集系统完整项目,面向需要构建爬虫或数据汇聚功能的Java开发者及大数据入门学习者。系统涵盖从HTTP接口抓取、多线程并发调度、Jsoup页面解析到数据清洗与JDBC存储的完整链路,适用于舆情监控、信息聚合、业务数据同步等典型场景。压缩包内共169个文件,包括35个Java源码、35个class字节码、39个jar依赖库、20个JSP页面以及XML和properties配置等,整体约22.54MB,项目结构清晰,便于按模块阅读与二次开发。目前已有332人学习下载,其中包含问卷调研、数据采集等业务模块的完整实现,展示了实际Web项目的分层设计。通过学习该项目,可掌握Java网络编程、并发采集、页面解析、持久化存储、日志监控等核心技能,是一份适合进阶与实战参考的完整代码包。

1. 数据采集系统的常见面貌

“java实现的数据采集系统.zip”初看是一个打包好的 Java 项目,实际上代表了一条把源数据搬到目标存储的工程链路:一次性导入、增量拉取、文件采集、接口轮询、日志和消息队列消费,最终都会落到“拉取-缓冲-清洗-写入”这四个动作上。常见误区是把采集系统想成一个大而全的平台,上来就画一大堆组件图,其实先搭一个可运行的骨架比什么都重要。骨架定下来之后,再围绕连接、并发、去重、重试和监控去加细节,系统才能从“能跑”走向“能用”。这套东西适合 Java 后端工程师用于对接第三方接口、处理文件交换、同步订单、聚合日志、迁移历史数据等场景,也是从 CRUD 开发走进数据方向最容易上手、也最容易讲清楚的一类项目。

2. 从最小可运行骨架开始:用原生 Java 搭出采集到入库的链路

2.1 阻塞队列解耦“采集”和“入库”两个动作

先不引入 Spring,也不引入 Flink 这类重框架。最常见的起点是一个采集线程、一个入库线程、一个阻塞队列:采集线程只负责任务拉取,入库线程只负责写库,中间用BlockingQueue传递原始数据。这样代码量小,每一条链路都能单独测试。

2.1.1 代码骨架
import java.sql.*; import java.time.LocalDateTime; import java.util.concurrent.*; public class PollingCollector { private static final BlockingQueue<String> RAW_QUEUE = new ArrayBlockingQueue<>(1000); public static void main(String[] args) throws Exception { // 采集线程:从数据源拉取原始数据 ExecutorService collectPool = Executors.newFixedThreadPool(3); // 入库线程:把队列中的数据批量写入数据库 ExecutorService storePool = Executors.newFixedThreadPool(2); for (int i = 0; i < 3; i++) { collectPool.submit(new PollTask(RAW_QUEUE)); } for (int i = 0; i < 2; i++) { storePool.submit(new StoreTask(RAW_QUEUE)); } collectPool.shutdown(); storePool.shutdown(); // 这里的 shutdown 只是不再接受新任务,正在执行的任务仍然会跑完 } }

采集任务的实现重点是“循环拉取”和“队列放入”:

class PollTask implements Runnable { private final BlockingQueue<String> queue; PollTask(BlockingQueue<String> queue) { this.queue = queue; } @Override public void run() { while (!Thread.currentThread().isInterrupted()) { try { String raw = fetchFromRemote(); if (raw != null) { queue.put(raw); // 队列满时会阻塞,产生背压 } } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } private String fetchFromRemote() { // 实际场景:调第三方接口、读文件、读DB,这里返回一个JSON字符串占位 return "{\"id\":1,\"updatedAt\":\"2024-01-01 23:01:01\"}"; } }

入库任务的实现重点是“从队列取数据”和“写库”:

class StoreTask implements Runnable { private final BlockingQueue<String> queue; StoreTask(BlockingQueue<String> queue) { this.queue = queue; } @Override public void run() { try (Connection conn = DriverManager.getConnection( "jdbc:mysql://localhost:3306/data_hub", "root", "root")) { while (!Thread.currentThread().isInterrupted()) { String record = queue.poll(2, TimeUnit.SECONDS); if (record != null) { insert(conn, record); } } } catch (Exception e) { // 必须打日志、重置连接,入库线程不能直接退出 } } private void insert(Connection conn, String record) throws SQLException { String sql = "INSERT INTO raw_log (content, created_at) VALUES (?, ?)"; try (PreparedStatement ps = conn.prepareStatement(sql)) { ps.setString(1, record); ps.setObject(2, LocalDateTime.now()); ps.executeUpdate(); } } }
2.1.2 这段代码里最关键的点是什么

ArrayBlockingQueue的容量是第一个要关注的参数。这里设成 1000,意味着当入库线程写入变慢时,采集线程最多只能在队列里塞 1000 条,再往下put会阻塞,这就是背压控制,比“无脑拉数据”安全得多。生产环境里我一般把队列容量放在配置项里,方便在监控里看到队列积压太多时临时调大,但真正的调优目标是让入库速度跟上采集速度,而不是无限放大队列。

采集线程数和入库线程数的配比要看 IO 类型。HTTP 调用属于高等待 IO,采集线程可以设置为 CPU 核数的 2 到 3 倍;JDBC 写入在事务提交时也会等待网络往返,同样需要多个连接并行。还是用上面这个例子,3 个采集线程加 2 个入库线程是一个低压力起点,后续用压测去调整。

2.2 文件目录监听:把 WatchService 接入同一个队列

接口轮询只是采集的一种来源。另一个常见场景是合作伙伴把数据文件丢到一个共享目录里,系统需要自动发现新文件并读取内容。Java 自带的WatchService可以作为目录监听的实现方案。

2.2.1 可复制的监听代码
import java.nio.file.*; public class FileWatcher { public static void watch(String dirPath) throws Exception { WatchService watchService = FileSystems.getDefault().newWatchService(); Paths.get(dirPath).register(watchService, StandardWatchEventKinds.ENTRY_CREATE, StandardWatchEventKinds.ENTRY_MODIFY); while (true) { WatchKey key = watchService.take(); // 阻塞等待文件事件 for (WatchEvent<?> event : key.pollEvents()) { Path fileName = (Path) event.context(); System.out.println("发现新文件: " + fileName); // 常见做法:把完整路径丢进一个线程池,由专门任务解析文件内容 } key.reset(); // 不 reset 会收不到后续事件 } } }

ENTRY_CREATE只能告诉你文件出现了,但如果上游是边写边传,文件可能还没写完。所以生产环境里更稳妥的做法是监听ENTRY_CREATE后,延迟几秒再去读,或者校验文件后缀和临时文件标志,例如.filepart代表传输未完成。ENTRY_MODIFY要慎用,一个持续写入的大文件会触发大量事件,容易打爆队列。

2.3 任务调度的选择:Timer、ScheduledExecutorService 还是 Quartz

对于定时轮询,Java 自带方案里有TimerScheduledExecutorServiceTimer的缺陷是单个线程执行任务,一个任务抛出未捕获异常会把整个调度器干掉,所以现在基本不用。ScheduledExecutorService可以指定多个线程,单个任务异常也不会影响到其他任务,是原生方案里的优先选择。只有需要 cron 表达式、分布式锁和失败转移时,再引入 Quartz。Quartz 的@Scheduled注解只是外表,底层依然是线程池加调度触发器,理解了这一点,面试时被问到“定时任务底层怎么实现”也不容易被带偏。

3. 三个必须显式配置的参数:批量大小、线程池容量和轮询间隔

3.1 批量大小由事务边界决定,不是越大越好

单条插入的性能一定不如批量插入,这是因为每次executeUpdate都有 SQL 解析、网络往返和事务提交的开销。常用的做法是把单条写改成addBatch方式,攒够固定条数后统一提交。注意:批量提交的条数直接决定单个事务的大小,批量太大反而会拖长锁持有时间,影响数据库并发。

3.1.1 参考的批量配置表
配置项建议初值适用场景调优方向
batchSize200~500MySQL 单表单写入观察 MySQL 响应时间,超过 500 时锁竞争会明显
flushInterval5000 ms低流量期兜底如果数据量很少,不可能无限等批次填满,定时强制提交
maxQueueSize1000~5000采集速度快于入库时先解决入库瓶颈,再决定是否调大
collectThreadsCPU 核数 × 2HTTP、RPC 调用为主看远端接口的响应时间和限流要求
storeThreadsCPU 核数 × 1JDBC 写库为主主要受数据库连接池大小限制

真实数据适合批量,但别忽略低峰期。设置flushInterval的意义在于,如果源端 10 分钟才来一条数据,你会一直等batchSize凑满,数据的落库延迟会被无限拉长。所以正确姿势是“优先攒批,超时强制提交”,两个条件满足一个就触发写入。

3.2 采集线程池的边界条件

Executors.newFixedThreadPool最简单,但从面试和排错角度,固定线程池的队列默认是无界的,意味着任务积压不会拒绝,只会把内存堆满。更成熟的写法是ThreadPoolExecutor显式指定核心线程数、最大线程数和拒绝策略:

ExecutorService collectPool = new ThreadPoolExecutor( 4, // 核心线程数 8, // 最大线程数,指核心线程之外的临时线程上限 60, TimeUnit.SECONDS, // 临时线程空闲回收时间 new ArrayBlockingQueue<>(500), // 任务队列 new ThreadPoolExecutor.CallerRunsPolicy() // 队列满时让提交线程自己执行 );

参数说明:核心线程数代表常态运行的任务数量,最大线程数代表突发流量下的最高任务数量,中间的空闲回收时间保证流量回落后线程数能降下来。拒绝策略用CallerRunsPolicy比直接抛异常安全,它表示队列满时后续任务由调用线程执行,在采集场景里相当于“采集线程自己等一等”,不会丢数据。AbortPolicy是默认策略,但在数据采集里我一般不直接抛异常,数据源通常不受控,宁可让采集线程阻塞。

3.3 轮询间隔的选择要看数据源水位

接口轮询间隔设置得过短,会给对方服务器造成无意义的压力,甚至触发限流;设置得过长,数据延迟又不可控。判断依据不应该是拍脑袋,而是源端数据更新频率。如果业务表每 5 分钟才产生一批新数据,轮询间隔设置成 10 秒没有意义,反而要关注每次拉取的是不是完整批次。

增量轮询时还要注意“窗口右移”的问题。pageNo翻页拉取数据时,如果源端一边查一边写入新数据,用 offset 分页容易出现重复或漏数据。常见做法是记录当前批次的最大主键 ID,作为下一批的起点:

SELECT * FROM source_order WHERE id > ? ORDER BY id ASC LIMIT 500;

这个写法依赖主键单调递增,对大多数订单、流水、日志类表都成立。如果源表的主键不是严格递增,就需要配合updated_at做时间窗口,并且引入游标延迟概念,都是增量采集里绕不开的细节。

4. 把数据变得可信:去重、幂等、失败重试与采集监控

4.1 去重:先想清楚“重复”的定义

数据采集的重复判断不能只靠数据库主键。同一份业务数据经过接口重试、文件重推、程序重启后,可能出现“内容相同但主键不同”的记录。在这种场景里,我会给每一条原始数据算一个指纹字段,用source_typesource_biz_idmd5(content)拼接后生成:

private String buildFingerprint(String sourceType, String bizId, String content) { String raw = sourceType + ":" + bizId + ":" + content; // 生产环境用 SHA-256 更稳妥,MD5 足够做去重 return DigestUtils.md5Hex(raw).toUpperCase(); }

指纹字段落库后必须建唯一索引,否则并发环境下两个任务同时插入同一条记录,应用层判断“不存在”和“插入”之间会产生竞态。唯一索引的存在让重复插入直接抛异常,再由INSERT ... ON DUPLICATE KEY UPDATE兜底,这是幂等落库的常见做法:

INSERT INTO raw_log (fingerprint, content, created_at) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE content = VALUES(content), updated_at = NOW();

参数说明:ON DUPLICATE KEY UPDATE在前半段插入失败时走后半段的更新分支。这样做不会让重复数据变成脏数据,而是把同一条指纹的记录更新时间刷新。对于数据采集场景,重复推送通常发生在程序重启后的新一轮拉取中,这个写法能避免大量“先查后写”逻辑。

4.2 失败重试:不能只靠 try-catch 吞异常

采集系统里最危险的操作是“捕获到异常后打印一行日志就继续循环”。这等同于把失败的责任甩给了下游,而下游数据库不会告诉你哪一批数据没进来。更合适的方案是区分两种失败:可重试的失败和不可重试的失败。

超时、数据库连接池耗尽、接口返回 5xx 是可重试的;参数错误、字段类型转换异常、数据库字段长度不足是不可重试的。可重试的数据要进入一个单独的重试队列,重试次数限制为 3 次,间隔采用指数退避:

public void retryWithBackoff(String record, int retryCount) { long waitMillis = (long) Math.pow(2, retryCount) * 1000; // 2s, 4s, 8s try { Thread.sleep(waitMillis); store(record); } catch (Exception e) { if (retryCount < 3) { retryWithBackoff(record, retryCount + 1); } else { // 推送到死信队列或重建一个待处理表,等人工排查 saveDeadLetter(record, e.getMessage()); } } }

这里幂等性往回又接上了第 4.1 节的唯一索引:因为重试意味着上一次执行可能已经写成功只是响应丢失,重试时再次插入就必须依赖唯一索引去重。没有这层兜底,指数退避的重试只会带来更多重复数据。

4.3 用少量代码实现采集水位监测

一个采集系统做得再花哨,最终要回答的问题是“数据延迟多久”和“积压了多少”。在纯 Java 项目里不需要一上来就接 Prometheus,可以先维护一组内存计数器:

public class CollectMetrics { private final AtomicLong pulledCount = new AtomicLong(0); // 已拉取条数 private final AtomicLong writtenCount = new AtomicLong(0); // 已写入条数 private final AtomicLong failedCount = new AtomicLong(0); // 失败条数 private volatile long lastPullTimestamp; // 最近一次拉取时间 public void onPull(long n) { pulledCount.addAndGet(n); lastPullTimestamp = System.currentTimeMillis(); } public void onWrite(long n) { writtenCount.addAndGet(n); } public void onFail(long n) { failedCount.addAndGet(n); } public long pending() { return pulledCount.get() - writtenCount.get(); } }

pending()方法返回的就是“拉取了但没写入”的积压量。当积压量持续增长时,说明消费端是瓶颈,需要先查数据库连接池、批量大小和慢 SQL。如果积压量一直在零附近,但业务反馈数据延迟大,就要看lastPullTimestamp距离当前时间多久,说明源端取数逻辑本身卡住了。这两类问题往往不在同一个排查路径上,分开监控比一个总指标更有用。

4.4 采集程序重启后如何保证不丢数据

基于内存队列的采集程序进程重启时,队列里还没入库的数据会直接丢失,这是在单机朴素实现里最容易暴露的短板。从轻量方案到重量方案依次是:落本地文件游标、写一份带状态的任务表、引入外部消息队列。

用游标文件记录“上次采集到哪”是最快的方案,每成功处理一批数据就更新一次游标值:

# sync-cursor.properties last.sync.time=2024-06-01 12:00:00 last.max.id=1029384

程序启动时读这个文件,接着游标位置继续拉。缺点是存在重复的可能,不能完全替代幂等去重。如果你们团队已经在用 Redis 或 MySQL,把游标存进去更合适,因为本地文件在容器化环境里会随 Pod 销毁而丢失。不管游标存哪里,重启流程一定是“先读状态,再启动采集”,顺序不能反。

5. 一个能验证全链路的实验:用增量轮询模拟订单同步

5.1 准备好源表和目标表

在本地 MySQL 建两张表,一张source_order模拟业务库订单表,一张sync_order模拟采集后的目标表。源表里必须有一个可比较的增量字段,最稳妥的是updated_at加主键id

CREATE TABLE source_order ( id BIGINT PRIMARY KEY AUTO_INCREMENT, order_no VARCHAR(32) NOT NULL, amount DECIMAL(10,2) NOT NULL, updated_at DATETIME NOT NULL ); CREATE TABLE sync_order ( id BIGINT PRIMARY KEY, order_no VARCHAR(32) NOT NULL, amount DECIMAL(10,2) NOT NULL, fingerprint VARCHAR(64) NOT NULL UNIQUE, synced_at DATETIME NOT NULL );

注意目标表里fingerprint字段加唯一索引,这是验证幂等性的前提。没有这个索引,后面的重复检查就演不了。

5.2 模拟业务方产生新数据

在源表里插入三条订单,模拟一个增量批次:

INSERT INTO source_order (order_no, amount, updated_at) VALUES ('ORD-2024001', 100.00, NOW()), ('ORD-2024002', 200.00, NOW()), ('ORD-2024003', 300.00, NOW());

然后执行采集程序的核心查询:

SELECT id, order_no, amount, updated_at FROM source_order WHERE updated_at > ? ORDER BY updated_at ASC, id ASC LIMIT 200;

这个 SQL 里ORDER BY updated_at ASC, id ASC要连在一起用。只按updated_at排序,当同一秒有大量并发写入时,上一批的最后一个游标值和下一批第一个游标值可能重叠,导致漏数据。加上id作为次级排序,并提供WHERE id > ?的双条件游标,能显著减少这个窗口。注意updated_at相等的场景仍然无法彻底避免重复,最终由目标表的唯一索引兜底,这也是为什么第 4 章先去重再谈同步顺序。

5.3 验证的三个关键数字

第一个验证点:运行一轮后,目标表的记录数不等于源表记录数时要能定位原因。跑一下两表对比查询:

SELECT COUNT(*) FROM source_order s LEFT JOIN sync_order t ON s.id = t.id WHERE t.id IS NULL;

返回 0 表示全部同步完成,返回非 0 表示有漏数据,需要看采集线程是否拉漏了某一页。

第二个验证点:把同一批数据手动重新插入源表并更新updated_at,再跑一遍采集,目标表不能出现重复订单。这段验证里fingerprint唯一索引就是守卫,命中重复时执行ON DUPLICATE KEY UPDATE只刷新synced_at

第三个验证点:在采集程序运行过程中直接结束进程,然后重新启动。启动后查日志确认游标是从上次记录的updated_atid恢复的,未处理完的数据会在新一轮轮询中被重新拉取。这个场景最容易暴露的问题是游标什么时候更新,它必须在你确认数据入库成功之后再写,写早了下一次启动就会跳过这些数据。

做完这三个验证,这套采集模块就已经具备了“可解释、可测试、可交接”的特质。再去扩展多数据源接入、并发调优或者引入消息队列,骨架都不会变。面试被问到数据采集项目时,能把上面两条 SQL 的游标条件和唯一索引的幂等作用讲透,比背一个完整的八股文框架图要更有说服力。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询