【发布时间】:2020-01-08 20:58:48
【问题描述】:
我正在使用 Azure Databricks 集群。
- Worker 和 Driver 类型:Standard_DS4_v2 规格:28.0 GB 内存,8 核心,1.5 DBU
- Databricks 运行时版本:6.1(包括 Apache Spark 2.4.4、Scala 2.11)
我有一个 pyspark 数据框 ttonag_algnd_prdictn_df1。这有大约 32,000 行。这个数据帧是使用 spark.read(...) 从 DB2 中提取的。我只使用 limit 关键字从中取出 10 行。
a = ttonag_algnd_prdictn_df1.limit(10)
a.show() 给出(为了便于阅读,我确实将文件放入文本文件并使其在 1 行中全部可见)
TONAG_ALGND_PRDICTN_ID,TONAG_MGT_YR,LINE_SGMT_NBR|TRAK_TYP_CD|BGN_MP_NBR|END_MP_NBR|TRAK_SDTRAK_NBR|ALGND_BGN_MP_NBR|ALGND_END_MP_NBR
1 2017 1 M 165.475 168.351 0 165.475 168.351
1 2018 1 M 165.475 168.351 0 165.475 168.351
1 2019 1 M 165.475 168.351 0 165.475 168.351
2 2016 1 M 395.225 405.698 0 395.225 405.698
2 2017 1 M 395.225 405.698 0 395.225 405.698
2 2018 1 M 395.225 405.698 0 395.225 405.698
2 2019 1 M 395.225 405.698 0 395.225 405.698
3 2016 1 M 412.005 422.198 0 412.005 422.198
3 2017 1 M 412.005 422.198 0 412.005 422.198
现在我做以下操作。
- 从 'a' 和 drop_duplicates 中选择列的子集。
unique_mp_pair_df = a.select("LINE_SGMT_NBR","TRAK_TYP_CD","TRAK_SDTRAK_NBR","ALGND_BGN_MP_NBR","ALGND_END_MP_NBR")
unique_mp_pair_df.show()
+-------------+-----------+---------------+----------------+----------------+
|LINE_SGMT_NBR|TRAK_TYP_CD|TRAK_SDTRAK_NBR|ALGND_BGN_MP_NBR|ALGND_END_MP_NBR|
+-------------+-----------+---------------+----------------+----------------+
| 1| M| 0 | 165.47500| 168.35100|
| 1| M| 0 | 165.47500| 168.35100|
| 1| M| 0 | 165.47500| 168.35100|
| 1| M| 0 | 165.47500| 168.35100|
| 1| M| 0 | 395.22500| 405.69800|
| 1| M| 0 | 395.22500| 405.69800|
| 1| M| 0 | 395.22500| 405.69800|
| 1| M| 0 | 395.22500| 405.69800|
| 1| M| 0 | 412.00500| 422.19800|
| 1| M| 0 | 412.00500| 422.19800|
+-------------+-----------+---------------+----------------+----------------+
unique_mp_pair_df = unique_mp_pair_df.drop_duplicates()
现在我希望这些行是唯一的。但是我得到的价值根本没有意义。
unique_mp_pair_df.show()
+-------------+-----------+---------------+----------------+----------------+
|LINE_SGMT_NBR|TRAK_TYP_CD|TRAK_SDTRAK_NBR|ALGND_BGN_MP_NBR|ALGND_END_MP_NBR|
+-------------+-----------+---------------+----------------+----------------+
| 7101| M| 0 | 11.29000| 24.88200|
+-------------+-----------+---------------+----------------+----------------+
上面是 ttonag_algnd_prdictn_df1 中的一行。但是在将行数限制为 10 之后,这不包括在内。如上所示,通过执行 a.show()
请有人帮助我理解这一点。我究竟做错了什么?非常感谢任何帮助。
【问题讨论】:
-
为什么会有反对票?我遵守了所有的规则。这是一个有效的问题。请帮助我了解它被否决的原因。
-
您能否检查 DB2 连接器中的 limit(n) 条记录是否发生谓词下推?
-
谓词下推是什么意思?我该怎么做?。
-
jaceklaskowski.gitbooks.io/mastering-spark-sql/… - 有关谓词下推的更多信息。如果可能的话,您能否使用解释方法分享数据框“a”的物理计划。
-
我使用了 a.cache()。解决了问题。感谢分享链接。会检查的
标签: apache-spark pyspark apache-spark-sql pyspark-sql databricks