Skip to content

Commit 1545ce3

Browse files
Fix columnar mismatch bug in iceberg dml when aqe enabled. (#13926)
1 parent ab1cda4 commit 1545ce3

10 files changed

Lines changed: 376 additions & 71 deletions

File tree

iceberg/src/main/scala/com/nvidia/spark/rapids/iceberg/IcebergProviderImpl.scala

Lines changed: 33 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,8 @@ package com.nvidia.spark.rapids.iceberg
1919
import scala.reflect.ClassTag
2020
import scala.util.Try
2121

22-
import com.nvidia.spark.rapids.{AppendDataExecMeta, AtomicCreateTableAsSelectExecMeta, AtomicReplaceTableAsSelectExecMeta, FileFormatChecks, GpuExec, GpuExpression, GpuRowToColumnarExec, GpuScan, IcebergFormatType, OverwriteByExpressionExecMeta, OverwritePartitionsDynamicExecMeta, RapidsConf, ScanMeta, ScanRule, ShimReflectionUtils, SparkPlanMeta, StaticInvokeMeta, TargetSize, WriteFileOp}
22+
import com.nvidia.spark.rapids.{AppendDataExecMeta, AtomicCreateTableAsSelectExecMeta, AtomicReplaceTableAsSelectExecMeta, FileFormatChecks, GpuExec, GpuExpression, GpuScan, IcebergFormatType, OverwriteByExpressionExecMeta, OverwritePartitionsDynamicExecMeta, RapidsConf, ScanMeta, ScanRule, ShimReflectionUtils, SparkPlanMeta, StaticInvokeMeta, WriteFileOp}
23+
import com.nvidia.spark.rapids.iceberg.IcebergProviderImpl.checkChildPlan
2324
import com.nvidia.spark.rapids.shims.{ReplaceDataExecMeta, WriteDeltaExecMeta}
2425
import org.apache.iceberg.spark.GpuTypeToSparkType.toSparkType
2526
import org.apache.iceberg.spark.functions._
@@ -31,6 +32,7 @@ import org.apache.spark.sql.catalyst.expressions.objects.StaticInvoke
3132
import org.apache.spark.sql.connector.read.Scan
3233
import org.apache.spark.sql.connector.write.Write
3334
import org.apache.spark.sql.execution.SparkPlan
35+
import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec
3436
import org.apache.spark.sql.execution.datasources.v2.{AppendDataExec, AtomicCreateTableAsSelectExec, AtomicReplaceTableAsSelectExec, GpuAppendDataExec, GpuOverwriteByExpressionExec, GpuOverwritePartitionsDynamicExec, GpuReplaceDataExec, GpuWriteDeltaExec, OverwriteByExpressionExec, OverwritePartitionsDynamicExec, ReplaceDataExec, WriteDeltaExec}
3537
import org.apache.spark.sql.execution.datasources.v2.rapids.{GpuAtomicCreateTableAsSelectExec, GpuAtomicReplaceTableAsSelectExec}
3638
import org.apache.spark.sql.types.{DateType, TimestampType}
@@ -177,6 +179,8 @@ class IcebergProviderImpl extends IcebergProvider {
177179
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)
178180

179181
GpuSparkWrite.tagForGpuCtas(cpuExec, meta)
182+
183+
checkChildPlan(meta)
180184
}
181185

182186
private def convertToGpu(
@@ -208,6 +212,8 @@ class IcebergProviderImpl extends IcebergProvider {
208212
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)
209213

210214
GpuSparkWrite.tagForGpuRtas(cpuExec, meta)
215+
216+
checkChildPlan(meta)
211217
}
212218

213219
private def convertToGpu(
@@ -238,15 +244,13 @@ class IcebergProviderImpl extends IcebergProvider {
238244
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)
239245

240246
GpuSparkWrite.tagForGpu(cpuExec.write, meta)
247+
248+
checkChildPlan(meta)
241249
}
242250

243251
private def convertToGpu(cpuExec: AppendDataExec, meta: AppendDataExecMeta): GpuExec = {
244-
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
245-
if (!child.supportsColumnar) {
246-
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
247-
}
248252
GpuAppendDataExec(
249-
child,
253+
meta.childPlans.head.convertIfNeeded(),
250254
cpuExec.refreshCache,
251255
GpuSparkWrite.convert(cpuExec.write))
252256
}
@@ -266,16 +270,14 @@ class IcebergProviderImpl extends IcebergProvider {
266270
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)
267271

268272
GpuSparkWrite.tagForGpu(cpuExec.write, meta)
273+
274+
checkChildPlan(meta)
269275
}
270276

271277
private def convertToGpu(cpuExec: OverwritePartitionsDynamicExec,
272278
meta: OverwritePartitionsDynamicExecMeta): GpuExec = {
273-
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
274-
if (!child.supportsColumnar) {
275-
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
276-
}
277279
GpuOverwritePartitionsDynamicExec(
278-
child,
280+
meta.childPlans.head.convertIfNeeded(),
279281
cpuExec.refreshCache,
280282
GpuSparkWrite.convert(cpuExec.write))
281283
}
@@ -295,16 +297,14 @@ class IcebergProviderImpl extends IcebergProvider {
295297
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)
296298

297299
GpuSparkWrite.tagForGpu(cpuExec.write, meta)
300+
301+
checkChildPlan(meta)
298302
}
299303

300304
private def convertToGpu(cpuExec: OverwriteByExpressionExec,
301305
meta: OverwriteByExpressionExecMeta): GpuExec = {
302-
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
303-
if (!child.supportsColumnar) {
304-
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
305-
}
306306
GpuOverwriteByExpressionExec(
307-
child,
307+
meta.childPlans.head.convertIfNeeded(),
308308
cpuExec.refreshCache,
309309
GpuSparkWrite.convert(cpuExec.write))
310310
}
@@ -366,15 +366,13 @@ class IcebergProviderImpl extends IcebergProvider {
366366
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)
367367

368368
GpuSparkWrite.tagForGpu(cpuExec.write, meta)
369+
370+
checkChildPlan(meta)
369371
}
370372

371373
private def convertToGpu(cpuExec: ReplaceDataExec, meta: ReplaceDataExecMeta): GpuExec = {
372-
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
373-
if (!child.supportsColumnar) {
374-
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
375-
}
376374
GpuReplaceDataExec(
377-
child,
375+
meta.childPlans.head.convertIfNeeded(),
378376
cpuExec.refreshCache,
379377
GpuSparkWrite.convert(cpuExec.write))
380378
}
@@ -395,17 +393,26 @@ class IcebergProviderImpl extends IcebergProvider {
395393
IcebergFormatType, WriteFileOp)
396394

397395
GpuSparkPositionDeltaWrite.tagForGpu(cpuExec.write, meta)
396+
397+
checkChildPlan(meta)
398398
}
399399

400400
private def convertToGpu(cpuExec: WriteDeltaExec, meta: WriteDeltaExecMeta): GpuExec = {
401-
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
402-
if (!child.supportsColumnar) {
403-
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
404-
}
405401
GpuWriteDeltaExec(
406-
child,
402+
meta.childPlans.head.convertIfNeeded(),
407403
cpuExec.refreshCache,
408404
cpuExec.projections,
409405
GpuSparkPositionDeltaWrite.convert(cpuExec.write))
410406
}
411407
}
408+
409+
object IcebergProviderImpl {
410+
def checkChildPlan[T <: SparkPlan](meta: SparkPlanMeta[T]): Unit = {
411+
if (meta.childPlans.nonEmpty) {
412+
val childMeta = meta.childPlans.head
413+
if (!childMeta.wrapped.isInstanceOf[AdaptiveSparkPlanExec] && !childMeta.canThisBeReplaced) {
414+
meta.willNotWorkOnGpu("Because child can't run gpu")
415+
}
416+
}
417+
}
418+
}

integration_tests/src/main/python/iceberg/iceberg_append_test.py

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ def test_insert_into_unpartitioned_table(spark_tmp_table_factory):
6363

6464
@iceberg
6565
@ignore_order(local=True)
66+
@allow_non_gpu('AppendDataExec')
6667
@pytest.mark.parametrize("partition_table", [True, False], ids=lambda x: f"partition_table={x}")
6768
def test_insert_into_unpartitioned_table_values(spark_tmp_table_factory,
6869
partition_table):
@@ -256,3 +257,47 @@ def insert_data(spark, table_name: str):
256257
conf = updated_conf)
257258

258259

260+
@iceberg
261+
@ignore_order(local=True)
262+
@pytest.mark.parametrize("partition_col_sql", [
263+
pytest.param(None, id="unpartitioned"),
264+
pytest.param("year(_c9)", id="year_partition"),
265+
])
266+
def test_insert_into_aqe(spark_tmp_table_factory, partition_col_sql):
267+
"""
268+
Test INSERT INTO with AQE enabled.
269+
"""
270+
table_prop = {"format-version": "2"}
271+
272+
# Configuration with AQE enabled
273+
conf = copy_and_update(iceberg_write_enabled_conf, {
274+
"spark.sql.adaptive.enabled": "true",
275+
"spark.sql.adaptive.coalescePartitions.enabled": "true"
276+
})
277+
278+
base_table_name = get_full_table_name(spark_tmp_table_factory)
279+
cpu_table_name = f"{base_table_name}_cpu"
280+
gpu_table_name = f"{base_table_name}_gpu"
281+
282+
# Create tables
283+
with_cpu_session(lambda spark: create_iceberg_table(
284+
cpu_table_name, partition_col_sql, table_prop,
285+
lambda sp: gen_df(sp, list(zip(iceberg_base_table_cols, iceberg_gens_list)))))
286+
with_cpu_session(lambda spark: create_iceberg_table(
287+
gpu_table_name, partition_col_sql, table_prop,
288+
lambda sp: gen_df(sp, list(zip(iceberg_base_table_cols, iceberg_gens_list)))))
289+
290+
# Insert data
291+
def insert_data(spark, table_name: str):
292+
df = gen_df(spark, list(zip(iceberg_base_table_cols, iceberg_gens_list)))
293+
view_name = spark_tmp_table_factory.get()
294+
df.createOrReplaceTempView(view_name)
295+
spark.sql(f"INSERT INTO {table_name} SELECT * FROM {view_name}")
296+
297+
with_gpu_session(lambda spark: insert_data(spark, gpu_table_name), conf=conf)
298+
with_cpu_session(lambda spark: insert_data(spark, cpu_table_name), conf=conf)
299+
300+
cpu_data = with_cpu_session(lambda spark: spark.table(cpu_table_name).collect())
301+
gpu_data = with_cpu_session(lambda spark: spark.table(gpu_table_name).collect())
302+
assert_equal_with_local_sort(cpu_data, gpu_data)
303+

integration_tests/src/main/python/iceberg/iceberg_ctas_test.py

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -255,6 +255,7 @@ def run_ctas(spark):
255255
@iceberg
256256
@ignore_order(local=True)
257257
@pytest.mark.parametrize("partition_table", [True, False], ids=lambda x: f"partition_table={x}")
258+
@allow_non_gpu('AtomicCreateTableAsSelectExec', 'AppendDataExec', 'ShuffleExchangeExec', 'SortExec', 'ProjectExec')
258259
def test_ctas_from_values(spark_tmp_table_factory,
259260
partition_table):
260261
table_prop = {
@@ -284,3 +285,40 @@ def execute_ctas_from_values(spark, target_table: str):
284285
cpu_data = with_cpu_session(lambda spark: spark.table(cpu_table).collect())
285286
gpu_data = with_cpu_session(lambda spark: spark.table(gpu_table).collect())
286287
assert_equal_with_local_sort(cpu_data, gpu_data)
288+
289+
290+
@iceberg
291+
@ignore_order(local=True)
292+
@pytest.mark.parametrize("partition_col_sql", [
293+
pytest.param(None, id="unpartitioned"),
294+
pytest.param("year(_c9)", id="triple_datetime_transforms"),
295+
])
296+
def test_ctas_aqe(spark_tmp_table_factory, partition_col_sql):
297+
"""
298+
Test CTAS with multiple partition transforms on the same column with AQE enabled.
299+
300+
This test reproduces NVBUGS-5689547 where the error "ROW BASED PROCESSING IS NOT SUPPORTED"
301+
occurs when writing to Iceberg tables with multiple partition transforms when AQE is enabled.
302+
303+
The issue manifests when:
304+
- AQE is enabled
305+
- Multiple partition transforms are applied (year, month, day, hour)
306+
- GpuShuffleCoalesceExec ends up as a child of GpuRowToColumnarExec
307+
"""
308+
table_prop = {
309+
"format-version": "2",
310+
}
311+
312+
df_gen = lambda spark: gen_df(spark, list(zip(iceberg_base_table_cols, iceberg_gens_list)))
313+
314+
# Configuration with AQE enabled (this is the key to reproducing the issue)
315+
conf = copy_and_update(iceberg_write_enabled_conf, {
316+
"spark.sql.adaptive.enabled": "true",
317+
"spark.sql.adaptive.coalescePartitions.enabled": "true"
318+
})
319+
320+
_assert_gpu_equals_cpu_ctas(spark_tmp_table_factory,
321+
df_gen,
322+
table_prop,
323+
partition_col_sql=partition_col_sql,
324+
conf=conf)

integration_tests/src/main/python/iceberg/iceberg_delete_test.py

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -460,3 +460,49 @@ def read_func(spark, table_name):
460460
"spark.rapids.sql.exec.WriteDeltaExec": "false"
461461
})
462462
)
463+
464+
465+
@allow_non_gpu("BatchScanExec", "ColumnarToRowExec")
466+
@iceberg
467+
@ignore_order(local=True)
468+
@pytest.mark.datagen_overrides(seed=DELETE_TEST_SEED, reason=DELETE_TEST_SEED_OVERRIDE_REASON)
469+
@pytest.mark.parametrize('update_mode', ['copy-on-write', 'merge-on-read'])
470+
@pytest.mark.parametrize("partition_col_sql", [
471+
pytest.param(None, id="unpartitioned"),
472+
pytest.param("year(_c9)", id="year_partition"),
473+
])
474+
def test_delete_aqe(spark_tmp_table_factory, update_mode, partition_col_sql):
475+
"""
476+
Test DELETE with AQE enabled.
477+
"""
478+
table_prop = {
479+
'format-version': '2',
480+
'write.delete.mode': update_mode
481+
}
482+
483+
# Configuration with AQE enabled
484+
conf = copy_and_update(iceberg_write_enabled_conf, {
485+
"spark.sql.adaptive.enabled": "true",
486+
"spark.sql.adaptive.coalescePartitions.enabled": "true"
487+
})
488+
489+
base_table_name = get_full_table_name(spark_tmp_table_factory)
490+
cpu_table = f"{base_table_name}_cpu"
491+
gpu_table = f"{base_table_name}_gpu"
492+
493+
def initialize_table(table_name):
494+
df_gen = lambda spark: gen_df(spark, list(zip(iceberg_base_table_cols, iceberg_gens_list)))
495+
create_iceberg_table(table_name, partition_col_sql, table_prop, df_gen)
496+
497+
with_cpu_session(lambda spark: initialize_table(cpu_table))
498+
with_cpu_session(lambda spark: initialize_table(gpu_table))
499+
500+
def delete_from_table(spark, table_name):
501+
spark.sql(f"DELETE FROM {table_name} WHERE _c2 % 3 = 0")
502+
503+
with_gpu_session(lambda spark: delete_from_table(spark, gpu_table), conf=conf)
504+
with_cpu_session(lambda spark: delete_from_table(spark, cpu_table), conf=conf)
505+
506+
cpu_data = with_cpu_session(lambda spark: spark.table(cpu_table).collect())
507+
gpu_data = with_cpu_session(lambda spark: spark.table(gpu_table).collect())
508+
assert_equal_with_local_sort(cpu_data, gpu_data)

0 commit comments

Comments
 (0)