apache / apache/carbondata

Why opened task less than available executors in case of insert into/load data

Open
#4,160 1 comment 0 reactions 0 assignees View on GitHub
Dominant language
Scala
Stars
1.5k
Forks
694
PR merge metrics
No merged PRs in 30d

Description

In case of insert into or load data, the total number of tasks in the stage is almost equal to the number of hosts, and in general it is much smaller than the available executors. The low parallelism of the stage results in slower execution. Why must the parallelism be constrained on the distinct host? Can start more tasks to increase parallelism and improve resource utilization? Thanks

org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala: loadDataFrame
```
/**
* Execute load process to load from input dataframe
*/
private def loadDataFrame(
sqlContext: SQLContext,
dataFrame: Option[DataFrame],
carbonLoadModel: CarbonLoadModel
): Array[(String, (LoadMetadataDetails, ExecutionErrors))] = {
try {
val rdd = dataFrame.get.rdd
val nodeNumOfData = rdd.partitions.flatMap[String, Array[String]] { p =>
DataLoadPartitionCoalescer.getPreferredLocs(rdd, p).map(_.host)
}.distinct.length
val nodes = DistributionUtil.ensureExecutorsByNumberAndGetNodeList(
nodeNumOfData,
sqlContext.sparkContext)
val newRdd = new DataLoadCoalescedRDD[Row](sqlContext.sparkSession, rdd, nodes.toArray
.distinct)

new NewDataFrameLoaderRDD(
sqlContext.sparkSession,
new DataLoadResultImpl(),
carbonLoadModel,
newRdd
).collect()
} catch {
case ex: Exception =>
LOGGER.error("load data frame failed", ex)
throw ex
}
}
```

Contributor guide

No contributing guide indexed for this repository

Research direction

Start with org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala and the loadDataFrame method shown in the issue. Trace how nodeNumOfData, ensureExecutorsByNumberAndGetNodeList, DataLoadCoalescedRDD, and NewDataFrameLoaderRDD determine task parallelism. Done means the load path can use available executor resources without the current host-based constraint, with execution behavior and resource utilization validated.

Written by the indexing model from the issue text.

Assessment

Tech stack
scala, spark
Domain
data-engineering, distributed-systems, performance
Issue type
Feature
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.