You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

从外部SQL数据库向Ignite导入千万级数据过慢问题排查

问题描述

我需要从外部SQL数据库向Ignite服务器导入1000万条数据,目前用Ignite缓存存储数据,已结合批处理与分页,通过Ignite的DataStreamer实现流式导入。但还未发送首批数据就收到「Possible too long JVM pause」告警,全量导入耗时约40分钟。

想请教:

  1. 我哪里操作有误?
  2. 除了增大堆内存外,还需做哪些调整?
  3. 是否存在更高效的外部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();
    }
}

问题分析与优化方案

一、当前操作的错误点

  1. 分页查询性能瓶颈:用Pageable生成的LIMIT offset, size语句,当offset增大时,数据库需要扫描大量无关数据才能返回目标页,导致查询耗时剧增,客户端线程阻塞,触发JVM停顿告警。
  2. DataStreamer重复创建:每次调用streamBulkData都新建IgniteDataStreamer实例,频繁销毁资源会带来额外开销,无法利用其批量缓冲优化。
  3. 线程与批大小不匹配:6个线程+10万条/批的配置,会导致客户端内存瞬间占用过高,触发GC;同时服务器端并行压力陡增,处理不过来。
  4. 手动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压力。

三、更高效的导入方法

  1. Ignite JDBC直接导入:使用Ignite JDBC驱动,让服务器直接从外部数据库拉取数据,执行INSERT INTO ignite_cache SELECT * FROM external_db_table,减少客户端中间环节。
  2. 服务器端分布式加载:通过IgniteCompute将数据加载任务提交到Ignite计算节点,让服务器直接从数据库拉取并写入缓存,避免客户端成为瓶颈。
  3. 文件中转导入:先将外部数据库数据导出为CSV/Parquet文件,再用Ignite的ignite.sh import命令或DataStreamer读取文件批量导入,降低数据库查询压力。

内容的提问来源于stack exchange,提问作者Dude Ramasamy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 19:04:57