ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Java数据采集系统实战:从阻塞队列到增量同步的完整链路

Java数据采集系统实战:从阻塞队列到增量同步的完整链路 简介这是一个基于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 BlockingQueueString 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 BlockingQueueString queue; PollTask(BlockingQueueString 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 BlockingQueueString queue; StoreTask(BlockingQueueString 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 自带方案里有Timer和ScheduledExecutorService。Timer的缺陷是单个线程执行任务一个任务抛出未捕获异常会把整个调度器干掉所以现在基本不用。ScheduledExecutorService可以指定多个线程单个任务异常也不会影响到其他任务是原生方案里的优先选择。只有需要 cron 表达式、分布式锁和失败转移时再引入 Quartz。Quartz 的Scheduled注解只是外表底层依然是线程池加调度触发器理解了这一点面试时被问到“定时任务底层怎么实现”也不容易被带偏。3. 三个必须显式配置的参数批量大小、线程池容量和轮询间隔3.1 批量大小由事务边界决定不是越大越好单条插入的性能一定不如批量插入这是因为每次executeUpdate都有 SQL 解析、网络往返和事务提交的开销。常用的做法是把单条写改成addBatch方式攒够固定条数后统一提交。注意批量提交的条数直接决定单个事务的大小批量太大反而会拖长锁持有时间影响数据库并发。3.1.1 参考的批量配置表配置项建议初值适用场景调优方向batchSize200500MySQL 单表单写入观察 MySQL 响应时间超过 500 时锁竞争会明显flushInterval5000 ms低流量期兜底如果数据量很少不可能无限等批次填满定时强制提交maxQueueSize10005000采集速度快于入库时先解决入库瓶颈再决定是否调大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_type、source_biz_id和md5(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.time2024-06-01 12:00:00 last.max.id1029384程序启动时读这个文件接着游标位置继续拉。缺点是存在重复的可能不能完全替代幂等去重。如果你们团队已经在用 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_at和id恢复的未处理完的数据会在新一轮轮询中被重新拉取。这个场景最容易暴露的问题是游标什么时候更新它必须在你确认数据入库成功之后再写写早了下一次启动就会跳过这些数据。做完这三个验证这套采集模块就已经具备了“可解释、可测试、可交接”的特质。再去扩展多数据源接入、并发调优或者引入消息队列骨架都不会变。面试被问到数据采集项目时能把上面两条 SQL 的游标条件和唯一索引的幂等作用讲透比背一个完整的八股文框架图要更有说服力。本文还有配套的精品资源点击获取
返回列表