Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
48 commits
Select commit Hold shift + click to select a range
f7f1723
Add test case for aqe
liurenjie1024 Dec 1, 2025
2569494
Fix test case
liurenjie1024 Dec 1, 2025
02ee6c7
Fix plan
liurenjie1024 Dec 1, 2025
f52e681
Columnar
liurenjie1024 Dec 1, 2025
080a40e
strategy rules
liurenjie1024 Dec 1, 2025
059fb44
Fix build break
liurenjie1024 Dec 1, 2025
6362f36
Add columnar write
liurenjie1024 Dec 1, 2025
cafc6ca
Register meta
liurenjie1024 Dec 1, 2025
4429085
transition
liurenjie1024 Dec 1, 2025
2df8f92
Fix row to c
liurenjie1024 Dec 2, 2025
0de977d
Fix build
liurenjie1024 Dec 2, 2025
77d993b
Fix build
liurenjie1024 Dec 2, 2025
5a2f7e5
Fix aqe
liurenjie1024 Dec 2, 2025
6de96ab
Order
liurenjie1024 Dec 2, 2025
aed4391
Fix build
liurenjie1024 Dec 2, 2025
bf8f6b2
Log
liurenjie1024 Dec 2, 2025
275d985
Fix aqe
liurenjie1024 Dec 2, 2025
d344d90
Print
liurenjie1024 Dec 2, 2025
227f5de
Columnar
liurenjie1024 Dec 2, 2025
9f4f62c
ctas aqe
liurenjie1024 Dec 2, 2025
d01556a
Fix all
liurenjie1024 Dec 2, 2025
0154e9c
Merge remote-tracking branch 'upstream/release/25.12' into ray/nvbugs…
liurenjie1024 Dec 2, 2025
f7abdb0
aqe for all
liurenjie1024 Dec 2, 2025
8668780
Project
liurenjie1024 Dec 2, 2025
774d690
rtas
liurenjie1024 Dec 2, 2025
1da27d0
Remove gpu r2c
liurenjie1024 Dec 2, 2025
65bbd02
Fix build
liurenjie1024 Dec 2, 2025
08f5faf
No project
liurenjie1024 Dec 2, 2025
0419f1f
Fix build
liurenjie1024 Dec 2, 2025
d0f5374
Restore unnecessary changes
liurenjie1024 Dec 2, 2025
c73e572
No gpu
liurenjie1024 Dec 3, 2025
4f80650
update aqe
liurenjie1024 Dec 3, 2025
3ccffb2
Skip ape
liurenjie1024 Dec 3, 2025
4360c96
fix aqe
liurenjie1024 Dec 3, 2025
6542cc5
cow only
liurenjie1024 Dec 3, 2025
3d79358
Fix aqe
liurenjie1024 Dec 3, 2025
df4451d
Fix aqe
liurenjie1024 Dec 3, 2025
7d9b755
Fix aqe
liurenjie1024 Dec 3, 2025
858a03d
Simplify some
liurenjie1024 Dec 3, 2025
c52b293
Fix merge more patterns
liurenjie1024 Dec 3, 2025
e2b1b29
fix ctas
liurenjie1024 Dec 4, 2025
c8363bb
fix fallback
liurenjie1024 Dec 4, 2025
425a806
Restore unnecessary changes
liurenjie1024 Dec 4, 2025
2a10445
Fix tests
liurenjie1024 Dec 4, 2025
79c4cd0
Fix test
liurenjie1024 Dec 4, 2025
49b25d7
Fix test
liurenjie1024 Dec 4, 2025
bd39a6f
Fix values test
liurenjie1024 Dec 4, 2025
c79085b
Fix comment
liurenjie1024 Dec 4, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,8 @@ package com.nvidia.spark.rapids.iceberg
import scala.reflect.ClassTag
import scala.util.Try

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

GpuSparkWrite.tagForGpuCtas(cpuExec, meta)

checkChildPlan(meta)
}

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

GpuSparkWrite.tagForGpuRtas(cpuExec, meta)

checkChildPlan(meta)
}

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

GpuSparkWrite.tagForGpu(cpuExec.write, meta)

checkChildPlan(meta)
}

private def convertToGpu(cpuExec: AppendDataExec, meta: AppendDataExecMeta): GpuExec = {
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
if (!child.supportsColumnar) {
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
}
GpuAppendDataExec(
child,
meta.childPlans.head.convertIfNeeded(),
cpuExec.refreshCache,
GpuSparkWrite.convert(cpuExec.write))
}
Expand All @@ -266,16 +270,14 @@ class IcebergProviderImpl extends IcebergProvider {
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)

GpuSparkWrite.tagForGpu(cpuExec.write, meta)

checkChildPlan(meta)
}

private def convertToGpu(cpuExec: OverwritePartitionsDynamicExec,
meta: OverwritePartitionsDynamicExecMeta): GpuExec = {
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
if (!child.supportsColumnar) {
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
}
GpuOverwritePartitionsDynamicExec(
child,
meta.childPlans.head.convertIfNeeded(),
cpuExec.refreshCache,
GpuSparkWrite.convert(cpuExec.write))
}
Expand All @@ -295,16 +297,14 @@ class IcebergProviderImpl extends IcebergProvider {
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)

GpuSparkWrite.tagForGpu(cpuExec.write, meta)

checkChildPlan(meta)
}

private def convertToGpu(cpuExec: OverwriteByExpressionExec,
meta: OverwriteByExpressionExecMeta): GpuExec = {
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
if (!child.supportsColumnar) {
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
}
GpuOverwriteByExpressionExec(
child,
meta.childPlans.head.convertIfNeeded(),
cpuExec.refreshCache,
GpuSparkWrite.convert(cpuExec.write))
}
Expand Down Expand Up @@ -366,15 +366,13 @@ class IcebergProviderImpl extends IcebergProvider {
FileFormatChecks.tag(meta, cpuExec.query.schema, IcebergFormatType, WriteFileOp)

GpuSparkWrite.tagForGpu(cpuExec.write, meta)

checkChildPlan(meta)
}

private def convertToGpu(cpuExec: ReplaceDataExec, meta: ReplaceDataExecMeta): GpuExec = {
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
if (!child.supportsColumnar) {
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
}
GpuReplaceDataExec(
child,
meta.childPlans.head.convertIfNeeded(),
cpuExec.refreshCache,
GpuSparkWrite.convert(cpuExec.write))
}
Expand All @@ -395,17 +393,26 @@ class IcebergProviderImpl extends IcebergProvider {
IcebergFormatType, WriteFileOp)

GpuSparkPositionDeltaWrite.tagForGpu(cpuExec.write, meta)

checkChildPlan(meta)
}

private def convertToGpu(cpuExec: WriteDeltaExec, meta: WriteDeltaExecMeta): GpuExec = {
var child: SparkPlan = meta.childPlans.head.convertIfNeeded()
if (!child.supportsColumnar) {
child = GpuRowToColumnarExec(child, TargetSize(meta.conf.gpuTargetBatchSizeBytes))
}
GpuWriteDeltaExec(
child,
meta.childPlans.head.convertIfNeeded(),
cpuExec.refreshCache,
cpuExec.projections,
GpuSparkPositionDeltaWrite.convert(cpuExec.write))
}
}

object IcebergProviderImpl {
def checkChildPlan[T <: SparkPlan](meta: SparkPlanMeta[T]): Unit = {
if (meta.childPlans.nonEmpty) {
val childMeta = meta.childPlans.head
if (!childMeta.wrapped.isInstanceOf[AdaptiveSparkPlanExec] && !childMeta.canThisBeReplaced) {
meta.willNotWorkOnGpu("Because child can't run gpu")
}
}
}
}
45 changes: 45 additions & 0 deletions integration_tests/src/main/python/iceberg/iceberg_append_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ def test_insert_into_unpartitioned_table(spark_tmp_table_factory):

@iceberg
@ignore_order(local=True)
@allow_non_gpu('AppendDataExec')
@pytest.mark.parametrize("partition_table", [True, False], ids=lambda x: f"partition_table={x}")
def test_insert_into_unpartitioned_table_values(spark_tmp_table_factory,
partition_table):
Expand Down Expand Up @@ -256,3 +257,47 @@ def insert_data(spark, table_name: str):
conf = updated_conf)


@iceberg
@ignore_order(local=True)
@pytest.mark.parametrize("partition_col_sql", [
pytest.param(None, id="unpartitioned"),
pytest.param("year(_c9)", id="year_partition"),
])
def test_insert_into_aqe(spark_tmp_table_factory, partition_col_sql):
"""
Test INSERT INTO with AQE enabled.
"""
table_prop = {"format-version": "2"}

# Configuration with AQE enabled
conf = copy_and_update(iceberg_write_enabled_conf, {
"spark.sql.adaptive.enabled": "true",
"spark.sql.adaptive.coalescePartitions.enabled": "true"
})

base_table_name = get_full_table_name(spark_tmp_table_factory)
cpu_table_name = f"{base_table_name}_cpu"
gpu_table_name = f"{base_table_name}_gpu"

# Create tables
with_cpu_session(lambda spark: create_iceberg_table(
cpu_table_name, partition_col_sql, table_prop,
lambda sp: gen_df(sp, list(zip(iceberg_base_table_cols, iceberg_gens_list)))))
with_cpu_session(lambda spark: create_iceberg_table(
gpu_table_name, partition_col_sql, table_prop,
lambda sp: gen_df(sp, list(zip(iceberg_base_table_cols, iceberg_gens_list)))))

# Insert data
def insert_data(spark, table_name: str):
df = gen_df(spark, list(zip(iceberg_base_table_cols, iceberg_gens_list)))
view_name = spark_tmp_table_factory.get()
df.createOrReplaceTempView(view_name)
spark.sql(f"INSERT INTO {table_name} SELECT * FROM {view_name}")

with_gpu_session(lambda spark: insert_data(spark, gpu_table_name), conf=conf)
with_cpu_session(lambda spark: insert_data(spark, cpu_table_name), conf=conf)

cpu_data = with_cpu_session(lambda spark: spark.table(cpu_table_name).collect())
gpu_data = with_cpu_session(lambda spark: spark.table(gpu_table_name).collect())
assert_equal_with_local_sort(cpu_data, gpu_data)

38 changes: 38 additions & 0 deletions integration_tests/src/main/python/iceberg/iceberg_ctas_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -255,6 +255,7 @@ def run_ctas(spark):
@iceberg
@ignore_order(local=True)
@pytest.mark.parametrize("partition_table", [True, False], ids=lambda x: f"partition_table={x}")
@allow_non_gpu('AtomicCreateTableAsSelectExec', 'AppendDataExec', 'ShuffleExchangeExec', 'SortExec', 'ProjectExec')
def test_ctas_from_values(spark_tmp_table_factory,
partition_table):
table_prop = {
Expand Down Expand Up @@ -284,3 +285,40 @@ def execute_ctas_from_values(spark, target_table: str):
cpu_data = with_cpu_session(lambda spark: spark.table(cpu_table).collect())
gpu_data = with_cpu_session(lambda spark: spark.table(gpu_table).collect())
assert_equal_with_local_sort(cpu_data, gpu_data)


@iceberg
@ignore_order(local=True)
@pytest.mark.parametrize("partition_col_sql", [
pytest.param(None, id="unpartitioned"),
pytest.param("year(_c9)", id="triple_datetime_transforms"),
])
def test_ctas_aqe(spark_tmp_table_factory, partition_col_sql):
"""
Test CTAS with multiple partition transforms on the same column with AQE enabled.

This test reproduces NVBUGS-5689547 where the error "ROW BASED PROCESSING IS NOT SUPPORTED"
occurs when writing to Iceberg tables with multiple partition transforms when AQE is enabled.

The issue manifests when:
- AQE is enabled
- Multiple partition transforms are applied (year, month, day, hour)
- GpuShuffleCoalesceExec ends up as a child of GpuRowToColumnarExec
"""
table_prop = {
"format-version": "2",
}

df_gen = lambda spark: gen_df(spark, list(zip(iceberg_base_table_cols, iceberg_gens_list)))

# Configuration with AQE enabled (this is the key to reproducing the issue)
conf = copy_and_update(iceberg_write_enabled_conf, {
"spark.sql.adaptive.enabled": "true",
"spark.sql.adaptive.coalescePartitions.enabled": "true"
})

_assert_gpu_equals_cpu_ctas(spark_tmp_table_factory,
df_gen,
table_prop,
partition_col_sql=partition_col_sql,
conf=conf)
46 changes: 46 additions & 0 deletions integration_tests/src/main/python/iceberg/iceberg_delete_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -460,3 +460,49 @@ def read_func(spark, table_name):
"spark.rapids.sql.exec.WriteDeltaExec": "false"
})
)


@allow_non_gpu("BatchScanExec", "ColumnarToRowExec")
@iceberg
@ignore_order(local=True)
@pytest.mark.datagen_overrides(seed=DELETE_TEST_SEED, reason=DELETE_TEST_SEED_OVERRIDE_REASON)
@pytest.mark.parametrize('update_mode', ['copy-on-write', 'merge-on-read'])
@pytest.mark.parametrize("partition_col_sql", [
pytest.param(None, id="unpartitioned"),
pytest.param("year(_c9)", id="year_partition"),
])
def test_delete_aqe(spark_tmp_table_factory, update_mode, partition_col_sql):
"""
Test DELETE with AQE enabled.
"""
table_prop = {
'format-version': '2',
'write.delete.mode': update_mode
}

# Configuration with AQE enabled
conf = copy_and_update(iceberg_write_enabled_conf, {
"spark.sql.adaptive.enabled": "true",
"spark.sql.adaptive.coalescePartitions.enabled": "true"
})

base_table_name = get_full_table_name(spark_tmp_table_factory)
cpu_table = f"{base_table_name}_cpu"
gpu_table = f"{base_table_name}_gpu"

def initialize_table(table_name):
df_gen = lambda spark: gen_df(spark, list(zip(iceberg_base_table_cols, iceberg_gens_list)))
create_iceberg_table(table_name, partition_col_sql, table_prop, df_gen)

with_cpu_session(lambda spark: initialize_table(cpu_table))
with_cpu_session(lambda spark: initialize_table(gpu_table))

def delete_from_table(spark, table_name):
spark.sql(f"DELETE FROM {table_name} WHERE _c2 % 3 = 0")

with_gpu_session(lambda spark: delete_from_table(spark, gpu_table), conf=conf)
with_cpu_session(lambda spark: delete_from_table(spark, cpu_table), conf=conf)

cpu_data = with_cpu_session(lambda spark: spark.table(cpu_table).collect())
gpu_data = with_cpu_session(lambda spark: spark.table(gpu_table).collect())
assert_equal_with_local_sort(cpu_data, gpu_data)
Loading