{"record":{"id":"d5e8eb5a95483f95","repo":"apache/iceberg","slug":"cannot-translate-spark-expression-sparkexpressio-d5e8eb","errorCode":null,"errorMessage":"Cannot translate Spark expression: $sparkExpression to data source filter","messagePattern":"Cannot translate Spark expression: \\$sparkExpression to data source filter","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"spark/v4.1/spark/src/main/scala/org/apache/spark/sql/execution/datasources/SparkExpressionConverter.scala","lineNumber":49,"sourceCode":"\nobject SparkExpressionConverter {\n\n  def convertToIcebergExpression(\n      sparkExpression: Expression): org.apache.iceberg.expressions.Expression = {\n    // Currently, it is a double conversion as we are converting Spark expression to Spark predicate\n    // and then converting Spark predicate to Iceberg expression.\n    // But these two conversions already exist and well tested. So, we are going with this approach.\n    DataSourceV2Strategy.translateFilterV2(sparkExpression) match {\n      case Some(filter) =>\n        val converted = SparkV2Filters.convert(filter)\n        if (converted == null) {\n          throw new IllegalArgumentException(\n            s\"Cannot convert Spark filter: $filter to Iceberg expression\")\n        }\n\n        converted\n      case _ =>\n        throw new IllegalArgumentException(\n          s\"Cannot translate Spark expression: $sparkExpression to data source filter\")\n    }\n  }\n\n  @throws[IcebergAnalysisException]\n  def collectResolvedSparkExpression(\n      session: SparkSession,\n      tableName: String,\n      where: String): Expression = {\n    val tableAttrs = session.table(tableName).queryExecution.analyzed.output\n    val unresolvedExpression = session.sessionState.sqlParser.parseExpression(where)\n    val filter = Filter(unresolvedExpression, DummyRelation(tableAttrs))\n    val optimizedLogicalPlan = session.sessionState.executePlan(filter).optimizedPlan\n    optimizedLogicalPlan\n      .collectFirst {\n        case filter: Filter => filter.condition\n        case _: DummyRelation => Literal.TrueLiteral\n        case _: LocalRelation => Literal.FalseLiteral","sourceCodeStart":31,"sourceCodeEnd":67,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/spark/v4.1/spark/src/main/scala/org/apache/spark/sql/execution/datasources/SparkExpressionConverter.scala#L31-L67","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["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.","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).","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.","As a last resort, extend the converter's match block with a case for the expression type and contribute a mapping to Iceberg expressions."],"exampleFix":"// before\ndf.filter(myUdf(col(\"ts\")) > 0) // UDF predicate cannot be translated\n// after\ndf.filter(col(\"ts\") > lit(0)) // built-in comparison is pushdown-convertible","handlingStrategy":"try-catch","validationCode":"// Before relying on pushdown, check the predicate is built from supported functions\nval unsupported = plan.collect { case f: Filter => f.condition }.exists(_.exists {\n  case _: ScalaUDF => true\n  case _ => false\n})\nif (unsupported) println(\"Predicate contains expressions that cannot be pushed to Iceberg\")","typeGuard":"def isPushable(e: org.apache.spark.sql.catalyst.expressions.Expression): Boolean =\n  !e.exists { case _: org.apache.spark.sql.catalyst.expressions.ScalaUDF => true; case _ => false }","tryCatchPattern":"try {\n  spark.table(\"catalog.db.t\").where(builtInPredicate).collect()\n} catch {\n  case e: IllegalArgumentException if e.getMessage.contains(\"Cannot translate Spark expression\") =>\n    // fall back to non-pushdown evaluation\n    df.filter(fallbackPredicate).collect()\n}","preventionTips":["Use only built-in comparisons, IN, AND/OR/NOT and null-safe equality in WHERE clauses intended for pushdown","Keep UDFs and complex expressions outside the pushed predicate","Keep the iceberg-spark artifact version aligned with your Spark version","Inspect the physical plan (df.explain) to confirm which filters were pushed"],"tags":["spark","filter-pushdown","unsupported-expression"],"backgroundTag":"unsupported-operation","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}