使用SparkSQL与es-hadoop从Elasticsearch读取全量文档仅获1万条问题
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.
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.sizetoo 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

