【发布时间】: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