如何用Spark并行读取文件夹中不同rowTag的XML文件并分结构生成DataFrame
Great question! Let's break this down into two clear parts: first, how to categorize your 1000 XML files into their three distinct structures, and second, how to parallelize the processing of each resulting DataFrame.
Step 1: Categorize XML Files by Structure
Since you can't distinguish file types by name, we first need to scan each file to identify its structure. We'll use lightweight XML parsing to avoid loading entire large files into memory, then map each file path to its structure type (a, b, or c).
1.1 Get All XML File Paths
First, retrieve the full path of every XML file in your target directory using Spark:
// Get paths of all XML files JavaRDD<String> xmlFilePaths = spark.sparkContext() .wholeTextFiles("/path/to/your/xml/folder") .map(Tuple2::_1); // Extract just the file path from the (path, content) tuple
1.2 Identify File Structure Type
Use a SAX parser (memory-efficient for large files) to check for unique tags that define each structure. For example, if your three structures use row tags a_record, b_record, c_record:
private static String determineStructureType(String filePath) { try (FileInputStream fis = new FileInputStream(new File(filePath))) { SAXParserFactory factory = SAXParserFactory.newInstance(); SAXParser parser = factory.newSAXParser(); final AtomicReference<String> structureType = new AtomicReference<>(); parser.parse(fis, new DefaultHandler() { @Override public void startElement(String uri, String localName, String qName, Attributes attributes) throws SAXException { // Stop parsing immediately once we find the defining tag switch(qName) { case "a_record": structureType.set("a"); throw new SAXException("Found structure a, halt parsing"); case "b_record": structureType.set("b"); throw new SAXException("Found structure b, halt parsing"); case "c_record": structureType.set("c"); throw new SAXException("Found structure c, halt parsing"); } } }); } catch (SAXException e) { // Catch our intentional halt exception and return the identified type if (structureType.get() != null) { return structureType.get(); } throw new RuntimeException("Failed to parse file " + filePath, e); } catch (Exception e) { throw new RuntimeException("Error accessing file " + filePath, e); } throw new IllegalArgumentException("Unknown structure for file: " + filePath); }
1.3 Group Files by Structure Type
Map each file path to its type, then collect the grouped results:
// Map files to their structure type JavaPairRDD<String, String> fileToStructure = xmlFilePaths.mapToPair(filePath -> new Tuple2<>(filePath, determineStructureType(filePath))); // Group files by structure type Map<String, List<String>> structureToFiles = fileToStructure.groupByKey() .mapValues(filePaths -> { List<String> pathList = new ArrayList<>(); filePaths.forEach(pathList::add); return pathList; }) .collectAsMap();
Step 2: Read Each Structure into a DataFrame
Now use Spark-XML to read each group of files into its own DataFrame, using the correct rowTag for each structure:
// Define rowTag mappings for each structure Map<String, String> structureToRowTag = new HashMap<>(); structureToRowTag.put("a", "a_record"); structureToRowTag.put("b", "b_record"); structureToRowTag.put("c", "c_record"); // Read files into DataFrames Map<String, Dataset<Row>> structureToDf = new HashMap<>(); for (Map.Entry<String, List<String>> entry : structureToFiles.entrySet()) { String structType = entry.getKey(); List<String> files = entry.getValue(); Dataset<Row> df = spark.read() .format("xml") .option("rowTag", structureToRowTag.get(structType)) .load(files.toArray(new String[0])); structureToDf.put(structType, df); }
Step 3: Parallelize Processing of the Three DataFrames
To avoid processing the DataFrames sequentially, use Java's concurrency utilities or Spark's asynchronous submission to run all three processing tasks in parallel.
Option 1: Use Java ExecutorService
This is straightforward for independent processing tasks:
// Create a thread pool sized to match the number of structures ExecutorService executor = Executors.newFixedThreadPool(3); // Define processing logic for each structure Runnable processA = () -> { Dataset<Row> dfA = structureToDf.get("a"); // Add your processing logic here (e.g., write to Parquet, transform data) dfA.write().format("parquet").save("/path/to/output/a"); }; Runnable processB = () -> { Dataset<Row> dfB = structureToDf.get("b"); dfB.write().format("parquet").save("/path/to/output/b"); }; Runnable processC = () -> { Dataset<Row> dfC = structureToDf.get("c"); dfC.write().format("parquet").save("/path/to/output/c"); }; // Submit tasks for parallel execution executor.submit(processA); executor.submit(processB); executor.submit(processC); // Wait for all tasks to complete executor.shutdown(); try { if (!executor.awaitTermination(1, TimeUnit.HOURS)) { executor.shutdownNow(); } } catch (InterruptedException e) { executor.shutdownNow(); }
Option 2: Use Spark's FutureAction
For tighter integration with Spark's context, use asynchronous job submission:
// Submit processing tasks asynchronously Future<Void> futureA = spark.sqlContext().submit(() -> { structureToDf.get("a").write().format("parquet").save("/path/to/output/a"); return null; }); Future<Void> futureB = spark.sqlContext().submit(() -> { structureToDf.get("b").write().format("parquet").save("/path/to/output/b"); return null; }); Future<Void> futureC = spark.sqlContext().submit(() -> { structureToDf.get("c").write().format("parquet").save("/path/to/output/c"); return null; }); // Wait for all jobs to finish futureA.get(); futureB.get(); futureC.get();
内容的提问来源于stack exchange,提问作者DINGJOY

