【问题标题】:Column name cannot be resolved in SparkSQL joinSparkSQL 连接中无法解析列名
【发布时间】:2016-10-04 06:07:40
【问题描述】:

我不确定为什么会这样。在 PySpark 中,我读入两个数据框并打印出它们的列名,它们与预期的一样,但是当执行 SQL 连接时,我收到一个错误,无法解析给定输入的列名。我已经简化了合并只是为了让它工作,但我需要添加更多的连接条件,这就是我使用 SQL 的原因(将添加:“and b.mnvr_bgn a.idx_trip_data”)。 df mnvr_temp_idx_prev_temp

中的“device_id”列似乎已重命名为“_col7”
mnvr_temp_idx_prev = mnvr_3.select('device_id', 'mnvr_bgn', 'mnvr_end')
print mnvr_temp_idx_prev.columns
['device_id', 'mnvr_bgn', 'mnvr_end']

raw_data_filtered = raw_data.select('device_id', 'trip_id', 'idx').groupby('device_id', 'trip_id').agg(F.max('idx').alias('idx_trip_end'))
print raw_data_filtered.columns
['device_id', 'trip_id', 'idx_trip_end']

raw_data_filtered.registerTempTable('raw_data_filtered_temp')
mnvr_temp_idx_prev.registerTempTable('mnvr_temp_idx_prev_temp') 
test = sqlContext.sql('SELECT a.device_id, a.idx_trip_end, b.mnvr_bgn, b.mnvr_end \
                          FROM raw_data_filtered_temp as a  \
                             INNER JOIN mnvr_temp_idx_prev_temp as b \
                                ON a.device_id = b.device_id')

回溯(最近一次调用最后一次):AnalysisException: u"cannot resolve 'b.device_id' given input columns: [_col7, trip_id, device_id, mnvr_end, mnvr_bgn, idx_trip_end]; line 1 pos 237"

感谢任何帮助!

【问题讨论】:

  • 请发布您的完整代码
  • 我的整个代码大约有 1000 行,所以这不是一个真正的选择
  • 您是否尝试过使用 DataFrames 进行 Join 而不是 sql 语句?没有太大区别,但想知道 Dataframes 中是否也出现同样的问题。

标签: pyspark apache-spark-sql pyspark-sql


【解决方案1】:

我建议至少在其中一个数据框中重命名“device_id”字段的名称。我稍微修改了您的查询并对其进行了测试(在 scala 中)。以下查询有效

test = sqlContext.sql("select * FROM raw_data_filtered_temp a INNER JOIN mnvr_temp_idx_prev_temp b ON a.device_id = b.device_id")
[device_id: string, mnvr_bgn: string, mnvr_end: string, device_id: string, trip_id: string, idx_trip_end: string]

现在,如果您在上述语句中执行“选择 *”,它将起作用。但是如果你尝试选择 'device_id',你会得到一个错误 "Reference 'device_id' is ambiguous" 。正如您在上面的“测试”数据框定义中看到的那样,它有两个具有相同名称(device_id)的字段。因此,为避免这种情况,我建议更改其中一个数据框中的字段名称。

mnvr_temp_idx_prev = mnvr_3.select('device_id', 'mnvr_bgn', 'mnvr_end')
                           .withColumnRenamned("device_id","device")  

raw_data_filtered = raw_data.select('device_id', 'trip_id', 'idx').groupby('device_id', 'trip_id').agg(F.max('idx').alias('idx_trip_end'))

现在使用数据框或 sqlContext

//using dataframes with multiple conditions
  val test = mnvr_temp_idx_prev.join(raw_data_filtered,$"device" === $"device_id"
                                                   && $"mnvr_bgn" < $"idx_trip_id","inner")

//在 SQL 上下文中

 test = sqlContext.sql("select * FROM raw_data_filtered_temp a INNER JOIN mnvr_temp_idx_prev_temp b ON a.device_id = b.device and a. idx_trip_id < b.mnvr_bgn")

以上查询可以解决您的问题。如果您的数据集太大,我建议不要在 Join 条件中使用 '>' 或 '

【讨论】:

  • 对于您的第一条评论,我确实尝试使用数据框连接并得到了同样的错误。重命名其中一个数据框中的列解决了这个问题!现在一切都按预期运行。谢谢!感谢您建议在 where 语句而不是 join 中使用 '>' 和 '
猜你喜欢
  • 1970-01-01
  • 2020-12-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-05-14
  • 2014-11-26
相关资源
最近更新 更多