{"record":{"id":"8dedb3a2d1bc8865","repo":"apache/iceberg","slug":"cannot-translate-spark-expression-sparkexpressio-8dedb3","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.2/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.2/spark/src/main/scala/org/apache/spark/sql/execution/datasources/SparkExpressionConverter.scala#L31-L67","documentation":"When a Spark expression cannot be translated into any DataSource V2 filter by Spark's `translateFilterV2` (returns None), Iceberg cannot represent it for pushdown and throws this IllegalArgumentException. It guards the first stage of the two-stage expression→filter→Iceberg-expression conversion.","triggerScenarios":"Calling `SparkExpressionConverter.convertToIcebergExpression` with a predicate Spark cannot lower to a V2 filter — e.g. subqueries, nondeterministic UDFs, unsupported expression types.","commonSituations":"Filters containing subqueries/EXISTS; deterministic-vs-nondeterministic UDF predicates; Spark version differences in translateFilterV2 coverage.","solutions":["Rewrite the filter with built-in operators supported by V2 pushdown","Materialize the subquery or join separately before filtering","Compute the predicate value first, then pass it as a literal","Keep unsupported predicates outside the pushed-down filter (post-scan filter)"],"exampleFix":"// before\n.where(\"col IN (SELECT id FROM other)\")\n// after\nval ids = spark.table(\"other\").collect().map(_.getInt(0)).toSeq\n.where(col(\"col\").isin(ids: _*))","handlingStrategy":"try-catch","validationCode":"// avoid subqueries/UDFs in pushed filters; precompute literals\nval threshold = spark.table(\"other\").agg(max(\"id\")).first().getInt(0)\nval df2 = df.filter(col(\"id\") < threshold)","typeGuard":null,"tryCatchPattern":"try { df.filter(expr) } catch { case e: IllegalArgumentException if e.getMessage.contains(\"Cannot translate Spark expression\") => /* rewrite or compute post-scan */ }","preventionTips":["Avoid subqueries and UDFs in filter expressions intended for pushdown","Precompute subquery results into literals or joins","Restrict filters to built-in operators","Verify with explain() that pushdown succeeds"],"tags":["spark","predicate-pushdown","expression-conversion"],"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"}