This is an automated email from the ASF dual-hosted git repository.
wForget pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new 1b545dfac8 [GLUTEN-12474][CORE] Preserve V1 write ordering for dynamic
partition writes (#12514)
1b545dfac8 is described below
commit 1b545dfac89cd5cd21e6d2ff621bdf957d39a74d
Author: Zhen Wang <[email protected]>
AuthorDate: Thu Jul 23 09:24:14 2026 +0800
[GLUTEN-12474][CORE] Preserve V1 write ordering for dynamic partition
writes (#12514)
* test GLUTEN-12474
* Preserve V1 write ordering when inserting local sorts
* add unit tests for other spark version
* fix
* fix spotless check
* Enhance local sort handling for GlutenPlan support
---
.../columnar/EnsureLocalSortRequirements.scala | 48 ++++++++++++++++++++--
.../datasources/GlutenV1WriteCommandSuite.scala | 35 ++++++++++++++++
.../datasources/GlutenV1WriteCommandSuite.scala | 35 ++++++++++++++++
.../datasources/GlutenV1WriteCommandSuite.scala | 35 ++++++++++++++++
.../datasources/GlutenV1WriteCommandSuite.scala | 35 ++++++++++++++++
.../org/apache/gluten/sql/shims/SparkShims.scala | 13 +++++-
.../gluten/sql/shims/spark34/Spark34Shims.scala | 15 +++++++
.../gluten/sql/shims/spark35/Spark35Shims.scala | 15 +++++++
.../gluten/sql/shims/spark40/Spark40Shims.scala | 15 +++++++
.../gluten/sql/shims/spark41/Spark41Shims.scala | 15 +++++++
10 files changed, 256 insertions(+), 5 deletions(-)
diff --git
a/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
b/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
index e17a8e7460..d22a71e7e9 100644
---
a/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
+++
b/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
@@ -16,11 +16,15 @@
*/
package org.apache.gluten.extension.columnar
+import org.apache.gluten.execution.GlutenPlan
import org.apache.gluten.extension.columnar.heuristic.HeuristicTransform
+import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.spark.sql.catalyst.expressions.SortOrder
import org.apache.spark.sql.catalyst.rules.Rule
-import org.apache.spark.sql.execution.{SortExec, SparkPlan}
+import org.apache.spark.sql.execution.{ColumnarWriteFilesExec, SortExec,
SparkPlan}
+import org.apache.spark.sql.execution.datasources.WriteFilesExec
+import org.apache.spark.sql.internal.SQLConf
/**
* This rule is similar with `EnsureRequirements` but only handle local
`SortExec`.
@@ -33,25 +37,61 @@ import org.apache.spark.sql.execution.{SortExec, SparkPlan}
object EnsureLocalSortRequirements extends Rule[SparkPlan] {
private lazy val transform: HeuristicTransform = HeuristicTransform.static()
+ private def numStaticPartitionCols(writeFiles: WriteFilesExec): Int = {
+ // HadoopFs writes include static partition columns in partitionColumns,
while Hive writes may
+ // only include the partition columns that are present in the write query.
+ val resolver = SQLConf.get.resolver
+ val staticPartitionNames = writeFiles.staticPartitions.keys
+ writeFiles.partitionColumns.takeWhile {
+ partitionColumn => staticPartitionNames.exists(resolver(_,
partitionColumn.name))
+ }.size
+ }
+
+ private def requiredChildOrdering(plan: SparkPlan): Seq[Seq[SortOrder]] = {
+ plan match {
+ // V1Writes assumes that the logical ordering it prepared is preserved
in the physical plan,
+ // so WriteFilesExec does not expose requiredChildOrdering itself.
Gluten may invalidate that
+ // ordering when it replaces a SortAggregateExec with a hash aggregate.
+ case writeFiles: WriteFilesExec
+ if ColumnarWriteFilesExec.OnNoopLeafPath.unapply(writeFiles).isEmpty
=>
+ Seq(
+ SparkShimLoader.getSparkShims.getV1WriteRequiredOrdering(
+ writeFiles.child.output,
+ writeFiles.partitionColumns,
+ writeFiles.bucketSpec,
+ writeFiles.options,
+ numStaticPartitionCols(writeFiles)))
+ case _ => plan.requiredChildOrdering
+ }
+ }
+
private def addLocalSort(
+ plan: SparkPlan,
originalChild: SparkPlan,
requiredOrdering: Seq[SortOrder]): SparkPlan = {
// FIXME: HeuristicTransform is costly. Re-applying it may cause
performance issues.
val newChild = SortExec(requiredOrdering, global = false, child =
originalChild)
- transform.apply(newChild)
+ (plan, originalChild) match {
+ case (_, child: GlutenPlan) if child.supportsColumnar =>
+ transform.apply(newChild)
+ case (parent: GlutenPlan, _) if parent.supportsColumnar =>
+ transform.apply(newChild)
+ case _ =>
+ newChild
+ }
}
override def apply(plan: SparkPlan): SparkPlan = {
plan.transformUp {
case p =>
- val newChildren = p.children.zip(p.requiredChildOrdering).map {
+ val newChildren = p.children.zip(requiredChildOrdering(p)).map {
case (child, requiredOrdering) =>
// If child.outputOrdering already satisfies the requiredOrdering,
// we do not need to sort.
if (SortOrder.orderingSatisfies(child.outputOrdering,
requiredOrdering)) {
child
} else {
- addLocalSort(child, requiredOrdering)
+ addLocalSort(p, child, requiredOrdering)
}
}
p.withNewChildren(newChildren)
diff --git
a/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
b/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index eb6794bba8..efd4105e68 100644
---
a/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++
b/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -98,6 +98,41 @@ class GlutenV1WriteCommandSuite
with GlutenV1WriteCommandSuiteBase
with GlutenSQLTestsBaseTrait {
+ testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+ withSQLConf(
+ "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+ "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+ withTable("gluten_12474_src", "gluten_12474_tgt") {
+ sql(
+ """
+ |CREATE TABLE gluten_12474_src USING ORC AS
+ |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+ | if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+ |FROM range(0, 10)
+ |""".stripMargin)
+
+ sql(
+ """
+ |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+ |USING ORC
+ |PARTITIONED BY (day STRING)
+ |""".stripMargin)
+
+ sql(
+ """
+ |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+ |SELECT k, max(v) AS m, day
+ |FROM gluten_12474_src
+ |GROUP BY day, k
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT k, m, day FROM gluten_12474_tgt"),
+ sql("SELECT k, v, day FROM gluten_12474_src"))
+ }
+ }
+ }
+
testGluten(
"SPARK-41914: v1 write with AQE and in-partition sorted - non-string
partition column") {
withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true") {
diff --git
a/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
b/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index 5fc887d8d4..6291ed3ac7 100644
---
a/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++
b/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -100,6 +100,41 @@ class GlutenV1WriteCommandSuite
with GlutenSQLTestsBaseTrait
with GlutenColumnarWriteTestSupport {
+ testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+ withSQLConf(
+ "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+ "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+ withTable("gluten_12474_src", "gluten_12474_tgt") {
+ sql(
+ """
+ |CREATE TABLE gluten_12474_src USING ORC AS
+ |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+ | if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+ |FROM range(0, 10)
+ |""".stripMargin)
+
+ sql(
+ """
+ |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+ |USING ORC
+ |PARTITIONED BY (day STRING)
+ |""".stripMargin)
+
+ sql(
+ """
+ |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+ |SELECT k, max(v) AS m, day
+ |FROM gluten_12474_src
+ |GROUP BY day, k
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT k, m, day FROM gluten_12474_tgt"),
+ sql("SELECT k, v, day FROM gluten_12474_src"))
+ }
+ }
+ }
+
testGluten(
"SPARK-41914: v1 write with AQE and in-partition sorted - non-string
partition column") {
withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true") {
diff --git
a/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
b/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index a287f5fffb..b8d9a1156e 100644
---
a/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++
b/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -100,6 +100,41 @@ class GlutenV1WriteCommandSuite
with GlutenSQLTestsBaseTrait
with GlutenColumnarWriteTestSupport {
+ testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+ withSQLConf(
+ "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+ "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+ withTable("gluten_12474_src", "gluten_12474_tgt") {
+ sql(
+ """
+ |CREATE TABLE gluten_12474_src USING ORC AS
+ |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+ | if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+ |FROM range(0, 10)
+ |""".stripMargin)
+
+ sql(
+ """
+ |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+ |USING ORC
+ |PARTITIONED BY (day STRING)
+ |""".stripMargin)
+
+ sql(
+ """
+ |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+ |SELECT k, max(v) AS m, day
+ |FROM gluten_12474_src
+ |GROUP BY day, k
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT k, m, day FROM gluten_12474_tgt"),
+ sql("SELECT k, v, day FROM gluten_12474_src"))
+ }
+ }
+ }
+
// TODO: fix in Spark-4.0
ignoreGluten(
"SPARK-41914: v1 write with AQE and in-partition sorted - non-string
partition column") {
diff --git
a/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
b/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index a287f5fffb..b8d9a1156e 100644
---
a/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++
b/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -100,6 +100,41 @@ class GlutenV1WriteCommandSuite
with GlutenSQLTestsBaseTrait
with GlutenColumnarWriteTestSupport {
+ testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+ withSQLConf(
+ "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+ "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+ withTable("gluten_12474_src", "gluten_12474_tgt") {
+ sql(
+ """
+ |CREATE TABLE gluten_12474_src USING ORC AS
+ |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+ | if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+ |FROM range(0, 10)
+ |""".stripMargin)
+
+ sql(
+ """
+ |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+ |USING ORC
+ |PARTITIONED BY (day STRING)
+ |""".stripMargin)
+
+ sql(
+ """
+ |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+ |SELECT k, max(v) AS m, day
+ |FROM gluten_12474_src
+ |GROUP BY day, k
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT k, m, day FROM gluten_12474_tgt"),
+ sql("SELECT k, v, day FROM gluten_12474_src"))
+ }
+ }
+ }
+
// TODO: fix in Spark-4.0
ignoreGluten(
"SPARK-41914: v1 write with AQE and in-partition sorted - non-string
partition column") {
diff --git
a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
index 750bea0c41..068ccdda6f 100644
--- a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
+++ b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
@@ -24,7 +24,8 @@ import org.apache.spark.broadcast.Broadcast
import org.apache.spark.internal.io.FileCommitProtocol
import org.apache.spark.sql.{AnalysisException, SparkSession}
import org.apache.spark.sql.catalyst.InternalRow
-import org.apache.spark.sql.catalyst.expressions.{Attribute, BinaryArithmetic,
Expression, InputFileBlockLength, InputFileBlockStart, InputFileName,
RaiseError, UnBase64}
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
+import org.apache.spark.sql.catalyst.expressions.{Attribute, BinaryArithmetic,
Expression, InputFileBlockLength, InputFileBlockStart, InputFileName,
RaiseError, SortOrder, UnBase64}
import org.apache.spark.sql.catalyst.plans.JoinType
import org.apache.spark.sql.catalyst.plans.QueryPlan
import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
@@ -129,6 +130,16 @@ trait SparkShims {
def enableNativeWriteFilesByDefault(): Boolean = false
+ // Planned V1 writes were introduced in Spark 3.4. Older versions do not
expose a required
+ // ordering utility and keep the default empty ordering.
+ // TODO: Remove this shim after dropping Spark 3.3 support.
+ def getV1WriteRequiredOrdering(
+ outputColumns: Seq[Attribute],
+ partitionColumns: Seq[Attribute],
+ bucketSpec: Option[BucketSpec],
+ options: Map[String, String],
+ numStaticPartitionCols: Int): Seq[SortOrder] = Seq.empty
+
def broadcastInternal[T: ClassTag](sc: SparkContext, value: T): Broadcast[T]
= {
// Since Spark 3.4, the `sc.broadcast` has been optimized to use
`sc.broadcastInternal`.
// More details see SPARK-39983.
diff --git
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
index 97ff19a84a..ae93d067a5 100644
---
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
+++
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
import org.apache.spark.sql.{AnalysisException, SparkSession}
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.analysis.DecimalPrecision
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
import org.apache.spark.sql.catalyst.expressions._
import org.apache.spark.sql.catalyst.expressions.aggregate._
import org.apache.spark.sql.catalyst.plans.QueryPlan
@@ -208,6 +209,20 @@ class Spark34Shims extends SparkShims {
override def enableNativeWriteFilesByDefault(): Boolean = true
+ override def getV1WriteRequiredOrdering(
+ outputColumns: Seq[Attribute],
+ partitionColumns: Seq[Attribute],
+ bucketSpec: Option[BucketSpec],
+ options: Map[String, String],
+ numStaticPartitionCols: Int): Seq[SortOrder] = {
+ V1WritesUtils.getSortOrder(
+ outputColumns,
+ partitionColumns,
+ bucketSpec,
+ options,
+ numStaticPartitionCols)
+ }
+
override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T):
Broadcast[T] = {
SparkContextUtils.broadcastInternal(sc, value)
}
diff --git
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
index d62fdfea19..08047c5757 100644
---
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
+++
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
import org.apache.spark.sql.{AnalysisException, SparkSession}
import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
import org.apache.spark.sql.catalyst.analysis.DecimalPrecision
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
import org.apache.spark.sql.catalyst.expressions._
import org.apache.spark.sql.catalyst.expressions.aggregate._
import org.apache.spark.sql.catalyst.plans.QueryPlan
@@ -249,6 +250,20 @@ class Spark35Shims extends SparkShims {
override def enableNativeWriteFilesByDefault(): Boolean = true
+ override def getV1WriteRequiredOrdering(
+ outputColumns: Seq[Attribute],
+ partitionColumns: Seq[Attribute],
+ bucketSpec: Option[BucketSpec],
+ options: Map[String, String],
+ numStaticPartitionCols: Int): Seq[SortOrder] = {
+ V1WritesUtils.getSortOrder(
+ outputColumns,
+ partitionColumns,
+ bucketSpec,
+ options,
+ numStaticPartitionCols)
+ }
+
override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T):
Broadcast[T] = {
SparkContextUtils.broadcastInternal(sc, value)
}
diff --git
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
index 5847e62c10..1e50984e19 100644
---
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
+++
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
import org.apache.spark.sql.{AnalysisException, SparkSession}
import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
import org.apache.spark.sql.catalyst.analysis.DecimalPrecisionTypeCoercion
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
import org.apache.spark.sql.catalyst.expressions._
import org.apache.spark.sql.catalyst.expressions.aggregate._
import org.apache.spark.sql.catalyst.plans.{JoinType, LeftSingle}
@@ -254,6 +255,20 @@ class Spark40Shims extends SparkShims {
override def enableNativeWriteFilesByDefault(): Boolean = true
+ override def getV1WriteRequiredOrdering(
+ outputColumns: Seq[Attribute],
+ partitionColumns: Seq[Attribute],
+ bucketSpec: Option[BucketSpec],
+ options: Map[String, String],
+ numStaticPartitionCols: Int): Seq[SortOrder] = {
+ V1WritesUtils.getSortOrder(
+ outputColumns,
+ partitionColumns,
+ bucketSpec,
+ options,
+ numStaticPartitionCols)
+ }
+
override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T):
Broadcast[T] = {
SparkContextUtils.broadcastInternal(sc, value)
}
diff --git
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
index bb0b94cc02..b0cd31be0e 100644
---
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
+++
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
import org.apache.spark.sql.{AnalysisException, SparkSession}
import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
import org.apache.spark.sql.catalyst.analysis.DecimalPrecisionTypeCoercion
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
import org.apache.spark.sql.catalyst.expressions._
import org.apache.spark.sql.catalyst.expressions.aggregate._
import org.apache.spark.sql.catalyst.plans.{JoinType, LeftSingle}
@@ -253,6 +254,20 @@ class Spark41Shims extends SparkShims {
override def enableNativeWriteFilesByDefault(): Boolean = true
+ override def getV1WriteRequiredOrdering(
+ outputColumns: Seq[Attribute],
+ partitionColumns: Seq[Attribute],
+ bucketSpec: Option[BucketSpec],
+ options: Map[String, String],
+ numStaticPartitionCols: Int): Seq[SortOrder] = {
+ V1WritesUtils.getSortOrder(
+ outputColumns,
+ partitionColumns,
+ bucketSpec,
+ options,
+ numStaticPartitionCols)
+ }
+
override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T):
Broadcast[T] = {
SparkContextUtils.broadcastInternal(sc, value)
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]