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

如何用Spark并行读取文件夹中不同rowTag的XML文件并分结构生成DataFrame

Solution: Separate XML Files by Structure & Parallel Processing with Spark-XML

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:02:47