基于Spark 1.6.0与CDH5.7.0,如何并行处理CSV中的表列表?
Hey there! Let's work through this parallel table processing problem you're hitting with Spark 1.6.0 on CDH 5.7.0.
First, let's break down why your current approaches might not be working:
- Using
collect()pulls all table names to the Driver node, so you're stuck processing tables one by one locally—this completely wastes Spark's distributed capabilities. - Scala Futures on their own (if run from the Driver) still tie processing to the Driver's resources, and you might not be properly leveraging Executor-side parallelism.
Here are two reliable, Spark-native solutions tailored to your environment:
Solution 1: Parallel Processing with mapPartitions + Local Thread Pools
This approach keeps your table list distributed, and uses thread pools on each Executor to process multiple tables in parallel within a single partition task.
import java.util.concurrent.Executors import scala.concurrent.{ExecutionContext, Future} import scala.util.{Success, Failure} // Initialize a thread pool per Executor (adjust the pool size based on your executor cores) val threadPoolSize = 4 // Match this to spark.executor.cores if possible val threadPool = Executors.newFixedThreadPool(threadPoolSize) implicit val ec = ExecutionContext.fromExecutor(threadPool) // Read your CSV table list as a DataFrame (keep it distributed—no collect!) val tableListDF = sqlContext.read .format("csv") .option("header", "true") .load("/path/to/your/tables.csv") // Process tables in parallel across partitions val processingResults = tableListDF.mapPartitions { tableRowIter => // For each table in the partition, kick off an async processing task val processingFutures = tableRowIter.map { row => val tableName = row.getString(0) // Adjust index based on your CSV structure Future { // Replace this with your actual table processing logic println(s"Processing table: $tableName on executor ${java.net.InetAddress.getLocalHost.getHostName}") val recordCount = sqlContext.sql(s"SELECT COUNT(*) FROM $tableName").first().getLong(0) (tableName, recordCount) } } // Wait for all futures in the partition to complete and return results processingFutures.map(_.value.get).iterator }.collect() // Output results processingResults.foreach { case (table, count) => println(s"Table $table processed successfully. Record count: $count") } // Clean up the thread pool threadPool.shutdown()
Why this works:
mapPartitionsruns on each Executor, so processing happens distributed across your cluster.- The local thread pool lets each Executor task handle multiple tables in parallel, maximizing resource utilization without overloading the Driver.
Solution 2: Using Spark's submitJob API
For more granular control over distributed tasks, you can use Spark's low-level submitJob API to parallelize table processing across partitions.
import org.apache.spark.SparkContext val sc: SparkContext = sqlContext.sparkContext // Read table names as an RDD (skip header if needed) val tableNamesRDD = sc.textFile("/path/to/your/tables.csv") .filter(line => !line.startsWith("table_name")) // Adjust header check to match your CSV // Define your table processing function def processSingleTable(tableName: String): String = { println(s"Processing table: $tableName") val count = sqlContext.sql(s"SELECT COUNT(*) FROM $tableName").first().getLong(0) s"$tableName -> Record count: $count" } // Submit parallel jobs across all partitions val finalResults = sc.submitJob[String, String, Array[String]]( rdd = tableNamesRDD, processFunc = (iter: Iterator[String]) => iter.map(processSingleTable).toArray, partitions = 0 until tableNamesRDD.partitions.length, // Process all partitions resultHandler = (partitionIndex: Int, results: Array[String]) => results, finalResultHandler = (allResults: Array[Array[String]]) => allResults.flatten ) // Print results finalResults.foreach(println)
Key Notes for Both Solutions:
- Avoid Driver-side bottlenecks: Never pull all table names to the Driver with
collect()unless your list is tiny. - Tune thread pool size: Match the pool size to your
spark.executor.coressetting to avoid resource contention on Executors. - Adjust partition count: If your CSV has a large number of tables, repartition it (e.g.,
tableListDF.repartition(10)) to distribute work across more Executors. - CDH 5.7.0 config tweaks: Ensure your Spark cluster has enough resources allocated—check
spark.executor.memoryandspark.executor.coresin Cloudera Manager to support parallel processing.
内容的提问来源于stack exchange,提问作者j pavan kumar
相关产品推荐
相关产品推荐

