apache/iceberg · error · IllegalArgumentException

Cannot translate Spark expression: $sparkExpression to data…

Error message

Cannot translate Spark expression: $sparkExpression to data source filter

What it means

Spark's filter-pushdown path failed to translate a resolved Spark Catalyst expression into an Iceberg data source filter. SparkExpressionConverter only knows how to map a fixed set of Spark predicates (comparisons, IN, AND/OR/NOT, etc.); anything outside that set, or a nested expression type it doesn't recognize, hits the default case and throws this IllegalArgumentException, aborting pushdown for that filter.

Solutions

  1. Simplify the WHERE clause: move unsupported predicates (UDFs, complex expressions) out of pushdown by wrapping them in a barrier, e.g. Spark's `spark.sql.optimizer.excludedRules` tuning or rewriting the query so only plain column comparisons reach the Iceberg scan.
  2. Check which exact expression fails (the message prints it) and rewrite it with supported functions (e.g. replace a UDF predicate with a built-in like col <=> value or standard comparisons).
  3. If the expression is a supported Spark predicate that still fails, verify your Iceberg Spark runtime version matches your Spark version; upgrade the iceberg-spark artifact.
  4. As a last resort, extend the converter's match block with a case for the expression type and contribute a mapping to Iceberg expressions.

Example fix

// before
df.filter(myUdf(col("ts")) > 0) // UDF predicate cannot be translated
// after
df.filter(col("ts") > lit(0)) // built-in comparison is pushdown-convertible
Defensive patterns

Strategy: try-catch

Validate before calling

// Before relying on pushdown, check the predicate is built from supported functions
val unsupported = plan.collect { case f: Filter => f.condition }.exists(_.exists {
  case _: ScalaUDF => true
  case _ => false
})
if (unsupported) println("Predicate contains expressions that cannot be pushed to Iceberg")

Type guard

def isPushable(e: org.apache.spark.sql.catalyst.expressions.Expression): Boolean =
  !e.exists { case _: org.apache.spark.sql.catalyst.expressions.ScalaUDF => true; case _ => false }

Try / catch

try {
  spark.table("catalog.db.t").where(builtInPredicate).collect()
} catch {
  case e: IllegalArgumentException if e.getMessage.contains("Cannot translate Spark expression") =>
    // fall back to non-pushdown evaluation
    df.filter(fallbackPredicate).collect()
}

Prevention

When it happens

Trigger: Calling SparkFilters.convert or SparkV2Filters.convert (or the extension-based filter pushdown in SparkTable/SparkScanBuilder) with a Catalyst expression not in the supported set, e.g. a user-defined function in WHERE, a struct/array comparison, an unsupported null-safe predicate variant, or a new Spark version introducing expression shapes the converter does not handle.

Common situations: Queries with UDF-based WHERE clauses pushed against an Iceberg table; upgrading Spark so the optimizer produces an expression shape the bundled converter predates; custom logical plans passing non-standard filter trees into V2 scan pushdown.

Understand the failure class

Background: UnsupportedOperationException and "is not supported" errors: when a library deliberately refuses a call — this error's family across 30 libraries.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/d5e8eb5a95483f95. Report an issue: GitHub.

Appendix: source

Thrown at spark/v4.1/spark/src/main/scala/org/apache/spark/sql/execution/datasources/SparkExpressionConverter.scala:49

object SparkExpressionConverter {

  def convertToIcebergExpression(
      sparkExpression: Expression): org.apache.iceberg.expressions.Expression = {
    // Currently, it is a double conversion as we are converting Spark expression to Spark predicate
    // and then converting Spark predicate to Iceberg expression.
    // But these two conversions already exist and well tested. So, we are going with this approach.
    DataSourceV2Strategy.translateFilterV2(sparkExpression) match {
      case Some(filter) =>
        val converted = SparkV2Filters.convert(filter)
        if (converted == null) {
          throw new IllegalArgumentException(
            s"Cannot convert Spark filter: $filter to Iceberg expression")
        }

        converted
      case _ =>
        throw new IllegalArgumentException(
          s"Cannot translate Spark expression: $sparkExpression to data source filter")
    }
  }

  @throws[IcebergAnalysisException]
  def collectResolvedSparkExpression(
      session: SparkSession,
      tableName: String,
      where: String): Expression = {
    val tableAttrs = session.table(tableName).queryExecution.analyzed.output
    val unresolvedExpression = session.sessionState.sqlParser.parseExpression(where)
    val filter = Filter(unresolvedExpression, DummyRelation(tableAttrs))
    val optimizedLogicalPlan = session.sessionState.executePlan(filter).optimizedPlan
    optimizedLogicalPlan
      .collectFirst {
        case filter: Filter => filter.condition
        case _: DummyRelation => Literal.TrueLiteral
        case _: LocalRelation => Literal.FalseLiteral

View on GitHub (pinned to 86d9c8fc54)