从外部SQL数据库向Ignite导入千万级数据过慢问题排查
问题描述
我需要从外部SQL数据库向Ignite服务器导入1000万条数据,目前用Ignite缓存存储数据,已结合批处理与分页,通过Ignite的DataStreamer实现流式导入。但还未发送首批数据就收到「Possible too long JVM pause」告警,全量导入耗时约40分钟。
想请教:
- 我哪里操作有误?
- 除了增大堆内存外,还需做哪些调整?
- 是否存在更高效的外部DB批量数据流式导入缓存的方法?
JVM配置(客户端与服务器端)
-Xms4g -Xmx4g -XX:+AlwaysPreTouch -XX:+UseG1GC -XX:+ScavengeBeforeFullGC -XX:ParallelGCThreads=8 -XX:ConcGCThreads=2 -XX:InitiatingHeapOccupancyPercent=55
相关代码
客户端批量数据流式处理代码
@Service public class Service { private static final Logger logger = LoggerFactory.getLogger(PosWavierService.class); private static final int BATCH_SIZE = 100_000; // 可调整批处理大小 private static final int NUM_THREADS = 6; private final ExecutorService executor = Executors.newFixedThreadPool(NUM_THREADS); @Autowired private IgniteCacheService igniteCacheService; @Autowired private Ignite ignite; @Autowired private Repo productRepo; public CompletableFuture<Void> processAllRecords(String cacheName) { long startTime = System.currentTimeMillis(); // 存储每个线程的CompletableFuture List<CompletableFuture<Void>> futures = new ArrayList<>(); // 并发提交数据获取与流式导入任务 AtomicInteger pageNumber = new AtomicInteger(0); for (int i = 0; i < NUM_THREADS; i++) { CompletableFuture<Void> future = CompletableFuture.supplyAsync(() -> { while (true) { List<ProductLines> records = fetchDataFromRepo(pageNumber.getAndIncrement(), BATCH_SIZE); if (records.isEmpty()) { break; } igniteCacheService.streamBulkData(cacheName, records); logger.info("已处理缓存{}的{}条记录", records.size(), cacheName); } return null; }); futures.add(future); } CompletableFuture<Void> combinedFuture = CompletableFuture.allOf( futures.toArray(new CompletableFuture[0])); combinedFuture.thenRun(() -> { long endTime = System.currentTimeMillis(); long totalTime = endTime - startTime; logger.info("处理所有记录总耗时:{}毫秒", totalTime); }); return combinedFuture; } }
数据库查询Repository代码
@Query(value = "SELECT * FROM table WHERE key IS NOT NULL AND key != '' AND key != ' ' ", nativeQuery = true) Page<Object> findRecordsWithPanNotNull(Pageable pageable);
streamBulkData函数代码
public void streamBulkData(String cacheName, List<Object> records) { try (IgniteDataStreamer<String, Object> streamer = ignite.dataStreamer(cacheName)){ // FileWriter writer = new FileWriter(KEYS_FILE_PATH, true)) { // 追加模式 streamer.allowOverwrite(true); streamer.perNodeBufferSize(1024); streamer.perNodeParallelOperations(8); streamer.skipStore(true); for (Product record : records) { String key = record.getPan_no(); if (key != null) { streamer.addData(key, record); //writer.write(key); } else { System.err.println("跳过空key的记录:" + record); } } streamer.flush(); } catch (CacheException e) { System.err.println("流式导入缓存出错:" + e.getMessage()); e.printStackTrace(); } }
问题分析与优化方案
一、当前操作的错误点
- 分页查询性能瓶颈:用
Pageable生成的LIMIT offset, size语句,当offset增大时,数据库需要扫描大量无关数据才能返回目标页,导致查询耗时剧增,客户端线程阻塞,触发JVM停顿告警。 - DataStreamer重复创建:每次调用
streamBulkData都新建IgniteDataStreamer实例,频繁销毁资源会带来额外开销,无法利用其批量缓冲优化。 - 线程与批大小不匹配:6个线程+10万条/批的配置,会导致客户端内存瞬间占用过高,触发GC;同时服务器端并行压力陡增,处理不过来。
- 手动flush打断缓冲机制:每个批次结束时调用
streamer.flush(),会强制中断DataStreamer的自动缓冲逻辑,降低批量导入效率。
二、除堆内存外的优化调整
1. 数据库查询优化
- 替换offset分页:改用基于主键/唯一键的范围查询,避免全表扫描。示例SQL:
SELECT id, pan_no, col1, col2 FROM table WHERE key IS NOT NULL AND key != '' AND key != ' ' AND id > ? ORDER BY id LIMIT ? - 减少查询字段:不要用
SELECT *,只查询需要导入的字段,降低数据传输量和内存占用。
2. DataStreamer配置优化
- 复用实例:全局复用1-2个
IgniteDataStreamer实例,避免重复创建销毁资源。 - 调整缓冲参数:
- 增大
perNodeBufferSize至8192或更高,积累更多数据后批量发送,减少网络交互。 - 调整
perNodeParallelOperations为服务器CPU核心数的1-2倍,避免过度并行引发资源竞争。 - 移除手动
flush(),让DataStreamer自动根据缓冲大小发送数据;若需要定期刷新,可设置streamer.autoFlushFrequency(1000)(单位:毫秒)。
- 增大
3. 客户端线程与批大小调整
- 减小单批次大小:将
BATCH_SIZE调整为2万-5万,避免单批次数据占用过多内存触发GC。 - 调整线程数:线程数不超过数据库连接池最大连接数,建议先尝试3-4个线程,观察服务器负载。
- 限流控制:添加简单的限流逻辑,避免短时间内发送大量数据压垮服务器。
4. JVM参数优化
- 调整G1GC参数:
- 增大
XX:MaxGCPauseMillis至500ms,减少GC频繁触发的可能。 - 根据CPU核心数调整
XX:ParallelGCThreads和XX:ConcGCThreads,比如8核CPU设置为6和2,避免GC线程占用过多CPU。
- 增大
- 启用
XX:+UseStringDeduplication:减少字符串对象内存占用,降低GC压力。
三、更高效的导入方法
- Ignite JDBC直接导入:使用Ignite JDBC驱动,让服务器直接从外部数据库拉取数据,执行
INSERT INTO ignite_cache SELECT * FROM external_db_table,减少客户端中间环节。 - 服务器端分布式加载:通过
IgniteCompute将数据加载任务提交到Ignite计算节点,让服务器直接从数据库拉取并写入缓存,避免客户端成为瓶颈。 - 文件中转导入:先将外部数据库数据导出为CSV/Parquet文件,再用Ignite的
ignite.sh import命令或DataStreamer读取文件批量导入,降低数据库查询压力。
内容的提问来源于stack exchange,提问作者Dude Ramasamy
相关产品推荐
相关产品推荐

