Files
sql-server-samples/samples/features/sql-big-data-cluster/spark/spark-sql.ipynb
T

11 KiB

Spark sample showing read/write methods

In this sample notebook, we will read CSV file from HDFS, write it as parquet file and save a Hive table definition. We will also run some Spark SQL commands using the Hive table.

In [1]:
# Read the clickstream CSV file(s) into a spark data frame, print schema & top rows
results = spark.read.option("inferSchema", "true").csv('/clickstream_data').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|
+-----------------+-----------------+------------+-----------+---------------+-----------+
|            38250|             4172|       67067|      11504|             40|      16819|
|            38251|            28919|       67090|      13782|             40|      11283|
|            38251|            77021|       67096|       4330|             40|      60107|
|            38251|            29023|       67109|      15796|             40|      31730|
|            38251|            54047|       67110|       9739|             40|      50449|
|            38251|            85733|       67117|       6843|             40|      39327|
|            38252|            53176|       67141|      12525|             40|      37913|
|            38252|            15873|       67153|      13008|             40|      92546|
|            38252|            39147|       67167|       5208|             40|      74534|
|            38252|            79540|       67171|      11552|             40|      94065|
|            38252|            35200|       67175|       9622|             40|      80502|
|            38253|            26068|       67191|       8585|             40|      43314|
|            38253|            63065|       67195|      17486|             40|      63793|
|            38253|             9687|       67214|       9856|             40|      92780|
|            38253|            18373|       67219|        406|             40|      38319|
|            38254|            80201|       67229|      13610|             40|      15342|
|            38254|            40058|       67239|      13594|             40|      41879|
|            38254|            79136|       67243|       1933|             40|      42095|
|            38254|            14684|       67244|      14267|             40|      39119|
|            38254|            36369|       67248|        641|             40|      82237|
+-----------------+-----------------+------------+-----------+---------------+-----------+
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
In [1]:
# Read the product reviews CSV files into a spark data frame, print schema & top rows
results = spark.read.option("inferSchema", "true").csv('/product_review_data').toDF(
            "pr_review_sk", "pr_review_content"
            )
results.printSchema()
results.show()
root
 |-- pr_review_sk: integer (nullable = true)
 |-- pr_review_content: string (nullable = true)

+------------+--------------------+
|pr_review_sk|   pr_review_content|
+------------+--------------------+
|       72621|Works fine. Easy ...|
|       89334|great product to ...|
|       89335|Next time will go...|
|       84259|Great Gift Great ...|
|       84398|After trip to Par...|
|       66434|Simply the best t...|
|       66501|This is the exact...|
|       66587|Not super magnet;...|
|       66680|Installed as bath...|
|       66694|Our home was buil...|
|       84489|Hi ;We are runnin...|
|       79052|Terra cotta is th...|
|       73034|One of my fingern...|
|       73298|We installed thes...|
|       66810|needed silicone c...|
|       66912|Great Gift Great ...|
|       67028|Laguiole knives a...|
|       89770|Good sound timers...|
|       84679|AWESOME FEEDBACK ...|
|       84953|love the retro gl...|
+------------+--------------------+
only showing top 20 rows
In [1]:
# Save results as parquet file and create hive table
results.write.format("parquet").mode("overwrite").saveAsTable("product_reviews")
In [1]:
# Execute Spark SQL commands
sqlDF = spark.sql("SELECT pr_review_sk, CHAR_LENGTH(pr_review_content) as len FROM product_reviews LIMIT 100")
sqlDF.show()
+------------+----+
|pr_review_sk| len|
+------------+----+
|       14868| 985|
|       14869|1601|
|       14875|1221|
|       14880| 665|
|       14886|  91|
|       14894| 697|
|       14899| 356|
|       14903|2361|
|       14908| 872|
|       14909|  74|
|       14917| 908|
|       14918|  50|
|       14919| 256|
|       14921| 723|
|       14925| 313|
|       14931|1304|
|       14939|1023|
|       14949| 552|
|       14954|2144|
|       14955| 123|
+------------+----+
only showing top 20 rows