mirror of
https://github.com/Microsoft/sql-server-samples.git
synced 2025-12-08 14:58:54 +00:00
8.3 KiB
8.3 KiB
In [1]:
# Read the CSV into a spark data frame, print schema & top rows
results = spark.read.option("inferSchema", "true").csv('/clickstream_data/web_clickstreams.csv').toDF(
"wcs_click_date_sk", "wcs_click_time_sk", "wcs_sales_sk", "wcs_item_sk", "wcs_web_page_sk", "wcs_user_sk"
)
results.printSchema()
results.show()root |-- wcs_click_date_sk: integer (nullable = true) |-- wcs_click_time_sk: integer (nullable = true) |-- wcs_sales_sk: integer (nullable = true) |-- wcs_item_sk: integer (nullable = true) |-- wcs_web_page_sk: integer (nullable = true) |-- wcs_user_sk: integer (nullable = true) +-----------------+-----------------+------------+-----------+---------------+-----------+ |wcs_click_date_sk|wcs_click_time_sk|wcs_sales_sk|wcs_item_sk|wcs_web_page_sk|wcs_user_sk| +-----------------+-----------------+------------+-----------+---------------+-----------+ | 36890| 40052| null| 4379| 34| null| | 36890| 41285| null| 6245| 34| null| | 36890| 23115| null| 13852| 34| null| | 36890| 17702| null| 15975| 34| null| | 36890| 62676| null| 2119| 34| null| | 36890| 34267| null| 10273| 34| null| | 36890| 8502| null| 17790| 34| null| | 36890| 54340| null| 3453| 34| null| | 36890| 54370| null| 6372| 34| null| | 36890| 6578| null| 17203| 34| null| | 36890| 75088| null| 4891| 34| null| | 36890| 23922| null| 11332| 34| null| | 36890| 28761| null| 4484| 34| null| | 36890| 21444| null| 5582| 34| null| | 36890| 58917| null| 8833| 34| null| | 36890| 27578| null| 8599| 34| null| | 36890| 8059| null| 6720| 34| null| | 36890| 43008| null| 17175| 34| null| | 36890| 4378| null| 10644| 34| null| | 36890| 55403| null| 8139| 34| null| +-----------------+-----------------+------------+-----------+---------------+-----------+ only showing top 20 rows
In [1]:
# Disable saving SUCCESS file
sc._jsc.hadoopConfiguration().set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")
# Print the current warehouse directory
print(spark.conf.get("spark.sql.warehouse.dir"))
# Save results as parquet file and create hive table
results.write.format("parquet").mode("overwrite").saveAsTable("web_clickstreams")
hdfs:///user/hive/warehouse
In [1]:
# Execute Spark SQL commands
sqlDF = spark.sql("SELECT * FROM web_clickstreams LIMIT 100")
sqlDF.show()
sqlDF = spark.sql("SELECT wcs_user_sk, COUNT(*)\
FROM web_clickstreams\
WHERE wcs_user_sk IS NOT NULL\
GROUP BY wcs_user_sk\
ORDER BY COUNT(*) DESC LIMIT 100")
sqlDF.show()+-----------------+-----------------+------------+-----------+---------------+-----------+ |wcs_click_date_sk|wcs_click_time_sk|wcs_sales_sk|wcs_item_sk|wcs_web_page_sk|wcs_user_sk| +-----------------+-----------------+------------+-----------+---------------+-----------+ | 37506| 7933| null| 1384| 2| 39437| | 37506| 56044| null| 14689| 2| 26419| | 37506| 52706| null| 8541| 2| 44016| | 37506| 67325| null| 16129| 2| 83371| | 37506| 84857| null| 1869| 2| 13090| | 37506| 49599| null| 2994| 2| 8940| | 37506| 78150| null| 11392| 2| 65633| | 37506| 38720| null| 14366| 2| 22281| | 37506| 79915| null| 11102| 2| 81755| | 37506| 67253| null| 5380| 2| 46868| | 37506| 6507| null| 6813| 2| 49363| | 37506| 18280| null| 1458| 2| 49363| | 37506| 72258| null| 2869| 2| 67756| | 37506| 8045| null| 615| 2| 86035| | 37506| 86164| null| 7000| 2| 94821| | 37506| 29724| null| 2767| 2| 94821| | 37506| 55471| null| 3584| 2| 62792| | 37506| 677| null| 1720| 2| 27212| | 37506| 66638| null| 9898| 2| 20370| | 37506| 48515| null| 9394| 2| 17157| +-----------------+-----------------+------------+-----------+---------------+-----------+ only showing top 20 rows +-----------+--------+ |wcs_user_sk|count(1)| +-----------+--------+ | 65042| 832| | 55928| 821| | 15570| 791| | 31138| 788| | 68188| 784| | 88205| 760| | 15678| 757| | 48063| 741| | 77518| 741| | 92978| 728| | 82129| 727| | 21700| 725| | 69707| 724| | 38895| 719| | 97643| 716| | 74426| 707| | 7813| 704| | 49528| 700| | 55766| 698| | 54355| 697| +-----------+--------+ only showing top 20 rows