【发布时间】:2020-01-03 14:29:19
【问题描述】:
我有两个数据集,我想通过 INNER JOIN 为我提供一个包含所需数据的全新表。我使用 SQL 并设法得到它。但是现在我想用map()和filter()试试,可以吗?
这是我使用 SPARK SQL 的代码:
import org.apache.spark.SparkContext
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
object hello {
def main(args: Array[String]): Unit = {
val conf = new SparkConf()
.setMaster("local")
.setAppName("quest9")
val sc = new SparkContext(conf)
val spark = SparkSession.builder().appName("quest9").master("local").getOrCreate()
val zip_codes = spark.read.format("csv").option("header", "true").load("/home/hdfs/Documents/quest_9/doc/zip.csv")
val census = spark.read.format("csv").option("header", "true").load("/home/hdfs/Documents/quest_9/doc/census.csv")
census.createOrReplaceTempView("census")
zip_codes.createOrReplaceTempView("zip")
//val query = spark.sql("SELECT * FROM census")
val query = spark.sql("SELECT DISTINCT census.Total_Males AS male, census.Total_Females AS female FROM census INNER JOIN zip ON census.Zip_Code=zip.Zip_Code WHERE zip.City = 'Inglewood' AND zip.County = 'Los Angeles'")
query.show()
query.write.parquet("/home/hdfs/Documents/population/census/IDE/census.parquet")
sc.stop()
}
}
【问题讨论】:
-
使用
dataframe.join()似乎比使用map或filter更明智。你为什么不使用它?见stackoverflow.com/questions/40343625/…或jaceklaskowski.gitbooks.io/mastering-spark-sql/…或stackoverflow.com/questions/36800174/… -
@GPI 因为我被指示这样做,出于某种原因
标签: scala apache-spark apache-spark-sql