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

使用SparkSQL与es-hadoop从Elasticsearch读取全量文档仅获1万条问题

Question: SparkSQL + es-hadoop only migrates 10000 documents from Elasticsearch index1 to index2

I'm trying to migrate all documents from Elasticsearch index1 to index2 using SparkSQL with the es-hadoop library, so I can perform intermediate data processing. I wrote the following Java class for this task:

Logger.getLogger("org.apache").setLevel(Level.WARN);
System.setProperty("hadoop.home.dir", "/opt/test");
SparkConf conf = new SparkConf()
    .setAppName("es-test")
    .setMaster("local[10]");
SparkSession sparkSession = SparkSession.builder()
    .config(conf)
    .getOrCreate();
Map<String, String> sourceTaskConfig = new HashMap<>();
sourceTaskConfig.put("es.nodes", "localhost");
sourceTaskConfig.put("es.port.size", "9200");
sourceTaskConfig.put("es.resource", "index1");
Dataset<Row> dataframe = JavaEsSparkSQL.esDF(new SQLContext(sparkSession), sourceTaskConfig);
//Here will be data processing
Map<String, String> destinationTaskConfig = new HashMap<>();
destinationTaskConfig.put("es.nodes", "localhost");
destinationTaskConfig.put("es.port.size", "9200");
JavaEsSparkSQL.saveToEs(dataframe, "index2", destinationTaskConfig);

I expect all data from index1 to be migrated to index2, but the task stops after only fetching 10000 documents. I've tried adjusting size limit parameters but still can't resolve this issue. Please tell me what's wrong with my implementation.


Answer

Hey there, let's figure out why you're only getting 10000 documents migrated. I've spotted a few key issues in your code and configuration that are almost certainly causing this problem:

1. Invalid Elasticsearch Port Configuration Key

You're using es.port.size which is not a valid es-hadoop configuration setting. The correct key for specifying the Elasticsearch port is es.port. This mistake could be causing unexpected connection behavior, which might interfere with the scroll pagination mechanism that es-hadoop uses to fetch large datasets.

2. Default Scroll Batch Size Limit

Under the hood, es-hadoop uses Elasticsearch's Scroll API to pull large datasets in batches. The default batch size for scroll is 10000—that's exactly why your job stops after retrieving 10000 documents. You need to override this default with the es.scroll.size parameter to fetch more records per batch.

3. Redundant SQLContext Creation

In Spark 2.x and later versions, SparkSession already encapsulates the SQLContext, so creating a new SQLContext(sparkSession) is unnecessary and could lead to unexpected context conflicts. You should use the SQLContext that's built into your SparkSession instead.

4. Missing Scroll Timeout

If your data processing takes significant time, the scroll context in Elasticsearch might expire before Spark finishes processing all batches. Adding es.scroll.timeout will extend the lifespan of the scroll context to prevent this.

Fixed Code Implementation

Here's the corrected version of your code with all these fixes applied, plus some additional optimizations for bulk writing:

import org.apache.log4j.Level;
import org.apache.log4j.Logger;
import org.apache.spark.SparkConf;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.elasticsearch.spark.sql.JavaEsSparkSQL;
import java.util.HashMap;
import java.util.Map;

public class EsMigration {
    public static void main(String[] args) {
        // Reduce verbose logging
        Logger.getLogger("org.apache").setLevel(Level.WARN);
        System.setProperty("hadoop.home.dir", "/opt/test");

        SparkConf conf = new SparkConf()
                .setAppName("es-test")
                .setMaster("local[10]");

        SparkSession sparkSession = SparkSession.builder()
                .config(conf)
                .getOrCreate();

        // Source ES configuration with fixes
        Map<String, String> sourceTaskConfig = new HashMap<>();
        sourceTaskConfig.put("es.nodes", "localhost");
        sourceTaskConfig.put("es.port", "9200"); // Fixed port configuration key
        sourceTaskConfig.put("es.resource", "index1");
        sourceTaskConfig.put("es.scroll.size", "50000"); // Increase batch size per scroll
        sourceTaskConfig.put("es.scroll.timeout", "10m"); // Extend scroll context lifespan
        sourceTaskConfig.put("es.query", "*"); // Explicitly fetch all documents (optional but safe)

        // Use SparkSession's built-in SQLContext
        Dataset<Row> dataframe = JavaEsSparkSQL.esDF(sparkSession.sqlContext(), sourceTaskConfig);

        // --- Your data processing logic here ---

        // Destination ES configuration with fixes and optimizations
        Map<String, String> destinationTaskConfig = new HashMap<>();
        destinationTaskConfig.put("es.nodes", "localhost");
        destinationTaskConfig.put("es.port", "9200"); // Fixed port configuration key
        destinationTaskConfig.put("es.batch.size.entries", "10000"); // Optimize bulk write batch size
        destinationTaskConfig.put("es.batch.write.refresh", "false"); // Disable refresh during bulk write for better performance
        destinationTaskConfig.put("es.mapping.id", "_id"); // Preserve original document IDs (optional if needed)

        // Save processed data to index2
        JavaEsSparkSQL.saveToEs(dataframe, "index2", destinationTaskConfig);

        // Clean up Spark resources
        sparkSession.stop();
    }
}

Additional Checks to Verify

  • Resource Limits: Ensure your Elasticsearch cluster has enough memory to handle larger scroll batches. Setting es.scroll.size too high could cause memory pressure on ES nodes—start with a reasonable value like 20000-50000 and adjust based on your cluster's capacity.
  • Data Filters: Confirm there are no implicit filters (like routing constraints or index aliases that only point to a subset of data) limiting the documents fetched. The es.query": "*" setting in the source config ensures you're pulling all documents.
  • Version Compatibility: Make sure your es-hadoop library version is compatible with both your Spark version and Elasticsearch version. Mismatched versions often lead to unexpected behavior like incomplete data pulls.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:21:43