【问题标题】:How to resolve the AnalysisException: resolved attribute(s) in Spark如何解决 AnalysisException:Spark 中已解析的属性
【发布时间】:2018-01-24 13:56:27
【问题描述】:
val rdd = sc.parallelize(Seq(("vskp", Array(2.0, 1.0, 2.1, 5.4)),("hyd",Array(1.5, 0.5, 0.9, 3.7)),("hyd", Array(1.5, 0.5, 0.9, 3.2)),("tvm", Array(8.0, 2.9, 9.1, 2.5))))
val df1= rdd.toDF("id", "vals")
val rdd1 = sc.parallelize(Seq(("vskp","ap"),("hyd","tel"),("bglr","kkt")))
val df2 = rdd1.toDF("id", "state")
val df3 = df1.join(df2,df1("id")===df2("id"),"left")

连接操作工作正常 但是当我重用 df2 时,我面临未解决的属性错误

val rdd2 = sc.parallelize(Seq(("vskp", "Y"),("hyd", "N"),("hyd", "N"),("tvm", "Y")))
val df4 = rdd2.toDF("id","existance")
val df5 = df4.join(df2,df4("id")===df2("id"),"left")

错误:org.apache.spark.sql.AnalysisException:已解析的属性 ID#426

【问题讨论】:

  • 这很可能与issues.apache.org/jira/browse/SPARK-10925有关,即id列的命名不明确
  • 但在第一种情况下它工作正常。我也提到了参考。我尝试将 df4 中的 id 重命名为 id_new。仍然无法解决错误。是因为 JAVARDD 的一些血统问题吗?我试着保留检查站。但还是同样的错误
  • 另见:stackoverflow.com/questions/40062298/… - 完整的错误消息是“已解决的属性缺失...”
  • 这可能是有史以来最糟糕/最骇人听闻的修复,但对数据帧进行别名化,即 df_alias = df.alias('df_alias) 并颠倒连接顺序,即将 df1_alias.join(df2_alias . . .) 更改为 df2_alias.join(df1_ailas . . .) 为我解决了这个问题跨度>

标签: java scala spark-dataframe


【解决方案1】:

@Json_Chans 的回答非常好,因为它不需要任何资源密集型操作。无论如何,在处理大量列时,您需要一些通用函数来动态处理这些内容,而不是手动编写数百个列。

幸运的是,您可以从 Dataframe 本身派生该函数,这样您就不需要任何额外的代码,除了单行代码(至少在 Python 和 pySpark 中):

import pyspark.sql.functions as f

df # Some Dataframe you have the "resolve(d) attribute(s)" error with

df = df.select([ f.col( column_name ).alias( column_name) for column_name in df.columns])

由于列的正确字符串表示仍存储在 Dataframe(df.columns: list) 的列属性中,您可以自行重置它 - 使用 .alias() 完成(注意:这仍然生成一个新的Dataframe,因为Dataframes 是不可变的,这意味着它们无法更改)。

【讨论】:

    【解决方案2】:

    在我的例子中,检查点原始数据框解决了这个问题。

    【讨论】:

      【解决方案3】:

      感谢Tomer's Answer

      对于 scala - 当我尝试使用自连接子句中的列时出现问题,使用方法修复它

      // To `and` all the column conditions
      def andAll(cols: Iterable[Column]): Column =
         if (cols.isEmpty) lit(true)
         else cols.tail.foldLeft(cols.head) { case (soFar, curr) => soFar.and(curr) }
      
      // To perform join different col name
      def renameColAndJoin(leftDf: DataFrame, joinCols: Seq[String], joinType: String = "inner")(rightDf: DataFrame): DataFrame = {
      
         val renamedCols: Seq[String]          = joinCols.map(colName => s"${colName}_renamed")
         val zippedCols: Seq[(String, String)] = joinCols.zip(renamedCols)
      
         val renamedRightDf: DataFrame = zippedCols.foldLeft(rightDf) {
           case (df, (origColName, renamedColName)) => df.withColumnRenamed(origColName, renamedColName)
         }
      
         val joinExpr: Column = andAll(zippedCols.map {
           case (origCol, renamedCol) => renamedRightDf(renamedCol).equalTo(rightDf(origCol))
         })
      
         leftDf.join(renamedRightDf, joinExpr, joinType)
      
      }
      

      【讨论】:

        【解决方案4】:

        这个问题真的浪费了我很多时间,我终于找到了一个简单的解决方案。

        在 PySpark 中,对于有问题的列,比如 colA,我们可以简单地使用

        import pyspark.sql.functions as F
        
        df = df.select(F.col("colA").alias("colA"))
        

        join 中使用df 之前。

        我认为这也适用于 Scala/Java Spark。

        【讨论】:

        • 谢天谢地,派生表只有两列,所以这个解决方法很容易编写,但真的很尴尬。这是与 2016 年相同的旧 SPARK-14948 错误吗,用于 Spark 1.6.0?
        • 哇!这对我有用,无法想象我浪费了这么多时间试图找出这个问题的根本原因:((
        • 非常聪明!谢谢!
        【解决方案5】:

        只需重命名您的列并输入相同的名称。 在 pyspark 中: 对于 df.columns 中的 i: df = df.withColumnRenamed(i,i)

        【讨论】:

        • 请解释一下并提供代码来帮助您回答
        • 我浏览了不同页面中的所有 cmets,但找不到答案。所以首先我尝试通过复制键和更改列名来加入表。所以我在表 a 中的 x 列对应于表 b 上的 y 列,并且错误消失了。然后我在聚合器函数之后重命名了两侧的列,并输入了与以前相同的名称。它像魔术一样工作。因此,如果在聚合函数之后有列 a、b、c,请将它们重命名为 a、b、c Ps。我已经提供了代码
        • 为我工作。感谢您提供这个简单但非常有效的解决方案。
        • 这应该更高!这些耗时问题的解决方法非常简单且有效。
        【解决方案6】:

        [TLDR]

        通过将中间DataFrame写入文件系统并再次读取来打破父DataFrame和派生DataFrame中列之间共享的AttributeReference

        例如:

        val df1 = spark.read.parquet("file1")
        df1.createOrReplaceTempView("df1")
        val df2 = spark.read.parquet("file2")
        df2.createOrReplaceTempView("df2")
        
        val df12 = spark.sql("""SELECT * FROM df1 as d1 JOIN df2 as d2 ON d1.a = d2.b""")
        df12.createOrReplaceTempView("df12")
        
        val df12_ = spark.sql(""" -- some transformation -- """)
        df12_.createOrReplaceTempView("df12_")
        
        val df3 = spark.read.parquet("file3")
        df3.createOrReplaceTempView("df3")
        
        val df123 = spark.sql("""SELECT * FROM df12_ as d12_ JOIN df3 as d3 ON d12_.a = d3.c""")
        df123.createOrReplaceTempView("df123")
        

        现在加入顶级 DataFrame 将导致“未解决的属性错误”

        val df1231 = spark.sql("""SELECT * FROM df123 as d123 JOIN df1 as d1 ON d123.a = d1.a""") 
        

        解决方案:d123.a 和 d1.a 共享相同的 AttributeReference 打破它 将中间表 df123 写入文件系统并再次读取。现在 df123write.a 和 d1.a 不共享 AttributeReference

        val df123 = spark.sql("""SELECT * FROM df12 as d12 JOIN df3 as d3 ON d12.a = d3.c""")
        df123.createOrReplaceTempView("df123")
        
        df123.write.parquet("df123.par")
        val df123write = spark.read.parquet("df123.par")
        spark.catalog.dropTempView("df123")
        df123write.createOrReplaceTempView("df123")
        
        val df1231 = spark.sql("""SELECT * FROM df123 as d123 JOIN df1 as d1 ON d123.a = d1.a""") 
        

        长篇大论:

        我们有复杂的 ETL,其中包含 DataFrame 的转换和自连接,在多个级别上执行。我们经常遇到“未解析的属性”错误,我们通过选择所需属性并在顶级表上执行连接而不是直接与顶级表连接来解决它,这暂时解决了问题,但是当我们对这些 DataFrame 应用更多转换并加入时对于任何顶级 DataFrame,“未解析的属性”错误再次引起了它的丑陋。

        发生这种情况是因为底层的 DataFrame 与派生它们的顶层 DataFrame 共享相同的 AttributeReference [more details]

        因此,我们通过仅写入 1 个中间转换的 DataFrame 并再次读取它并继续我们的 ETL 来打破这种引用共享。这打破了底部 DataFrame 和顶部 DataFrame 之间共享 AttributeReference 的问题,我们再也不会遇到“未解决的属性”错误。

        这对我们很有效,因为当我们从顶层 DataFrame 转移到底层执行转换并加入我们的数据比我们开始的初始 DataFrames 缩小时,它还提高了我们的性能,因为数据大小更小并且 spark 不必遍历DAG 一直到最后一个持久化的 DataFrame。

        【讨论】:

          【解决方案7】:

          在我的情况下,这个错误出现在同一张表的自连接期间。 我遇到了 Spark SQL 而不是数据框 API 的以下问题:

          org.apache.spark.sql.AnalysisException: Resolved attribute(s) originator#3084,program_duration#3086,originator_locale#3085 missing from program_duration#1525,guid#400,originator_locale#1524,EFFECTIVE_DATETIME_UTC#3157L,device_timezone#2366,content_rpd_id#734L,originator_sublocale#2355,program_air_datetime_utc#3155L,originator#1523,master_campaign#735,device_provider_id#2352 in operator !Deduplicate [guid#400, program_duration#3086, device_timezone#2366, originator_locale#3085, originator_sublocale#2355, master_campaign#735, EFFECTIVE_DATETIME_UTC#3157L, device_provider_id#2352, originator#3084, program_air_datetime_utc#3155L, content_rpd_id#734L]. Attribute(s) with the same name appear in the operation: originator,program_duration,originator_locale. Please check if the right attribute(s) are used.;;
          

          之前我使用以下查询,

              SELECT * FROM DataTable as aext
                       INNER JOIN AnotherDataTable LAO 
          ON aext.device_provider_id = LAO.device_provider_id 
          

          在加入之前只选择所需的列为我解决了这个问题。

                SELECT * FROM (
              select distinct EFFECTIVE_DATE,system,mso_Name,EFFECTIVE_DATETIME_UTC,content_rpd_id,device_provider_id 
          from DataTable 
          ) as aext
                   INNER JOIN AnotherDataTable LAO ON aext.device_provider_id = LAO.device_provider_id 
          

          【讨论】:

            【解决方案8】:

            根据我的经验,我们有 2 个解决方案 1) 克隆 DF 2) 在加入表之前重命名有歧义的列。 (不要忘记删除重复的连接键)

            我个人更喜欢第二种方法,因为在第一种方法中克隆 DF 需要时间,尤其是在数据量很大的情况下。

            【讨论】:

            • "第一种方法克隆 DF 需要时间,尤其是在数据量很大的情况下":丢弃克隆选项以解决大表问题
            【解决方案9】:

            如果您执行以下操作,它将起作用。

            假设你有一个数据框。 df1 如果你想交叉加入同一个数据框,你可以使用下面的

            df1.toDF("ColA","ColB").as("f_df").join(df1.toDF("ColA","ColB").as("t_df"), 
               $"f_df.pcmdty_id" === 
               $"t_df.assctd_pcmdty_id").select($"f_df.pcmdty_id",$"f_df.assctd_pcmdty_id")
            

            【讨论】:

              【解决方案10】:

              如果您有 df1,并且 df2 从 df1 派生,请尝试重命名 df2 中的所有列,以便在连接后没有两列具有相同的名称。所以在加入之前:

              所以而不是df1.join(df2...

              # Step 1 rename shared column names in df2.
              df2_renamed = df2.withColumnRenamed('columna', 'column_a_renamed').withColumnRenamed('columnb', 'column_b_renamed')
              
              # Step 2 do the join on the renamed df2 such that no two columns have same name.
              df1.join(df2_renamed)
              

              【讨论】:

              • 为我工作。谢谢
              • 很容易理解。谢谢!! :)
              • 这个完美的工作!优雅地解释。这节省了我的一天:)非常感谢:)
              【解决方案11】:

              尝试在两个连续的连接中使用一个 DataFrame 时,我遇到了同样的问题。

              问题是:DataFrame A 有 2 列(我们称它们为 x 和 y),DataFrame B 也有 2 列(我们称它们为 w 和 z)。我需要在 x=z 上将 A 与 B 连接起来,然后在 y=z 上将它们连接在一起。

              (A join B on A.x=B.z) as C join B on C.y=B.z
              

              我得到的确切错误是在第二次加入时它抱怨“resolved attribute(s) B.z#1234 ...”。

              根据@Erik 提供的链接以及其他一些博客和问题,我收集到我需要一个 B 的克隆。

              这是我所做的:

              val aDF = ...
              val bDF = ...
              val bCloned = spark.createDataFrame(bDF.rdd, bDF.schema)
              aDF.join(bDF, aDF("x") === bDF("z")).join(bCloned, aDF("y") === bCloned("z"))
              

              【讨论】:

                【解决方案12】:

                对于java开发者,请尝试调用此方法:

                private static Dataset<Row> cloneDataset(Dataset<Row> ds) {
                    List<Column> filterColumns = new ArrayList<>();
                    List<String> filterColumnsNames = new ArrayList<>();
                    scala.collection.Iterator<StructField> it = ds.exprEnc().schema().toIterator();
                    while (it.hasNext()) {
                        String columnName = it.next().name();
                        filterColumns.add(ds.col(columnName));
                        filterColumnsNames.add(columnName);
                    }
                    ds = ds.select(JavaConversions.asScalaBuffer(filterColumns).seq()).toDF(scala.collection.JavaConverters.asScalaIteratorConverter(filterColumnsNames.iterator()).asScala().toSeq());
                    return ds;
                }
                

                在加入之前的两个数据集上,它将数据集克隆到新数据集:

                df1 = cloneDataset(df1); 
                df2 = cloneDataset(df2);
                Dataset<Row> join = df1.join(df2, col("column_name"));
                // if it didn't work try this
                final Dataset<Row> join = cloneDataset(df1.join(df2, columns_seq)); 
                

                【讨论】:

                  【解决方案13】:

                  正如我在评论中提到的,它与https://issues.apache.org/jira/browse/SPARK-10925 相关,更具体地说与https://issues.apache.org/jira/browse/SPARK-14948 有关。重复使用引用会在命名中产生歧义,因此您必须克隆 df - 有关示例,请参见 https://issues.apache.org/jira/browse/SPARK-14948 中的最后一条评论。

                  【讨论】:

                  猜你喜欢
                  • 2021-08-13
                  • 1970-01-01
                  • 1970-01-01
                  • 1970-01-01
                  • 2020-02-01
                  • 1970-01-01
                  • 2012-11-02
                  • 1970-01-01
                  • 1970-01-01
                  相关资源
                  最近更新 更多