เริ่มต้นใช้งานเริ่มต้นใช้งานได้ฟรี

นำ Spark มาใช้กับข้อมูลของคุณ

ในแบบฝึกหัดที่แล้ว เราได้เรียนรู้วิธีย้ายข้อมูลจาก Spark ไปยัง pandas แต่บางครั้งอาจต้องการทำตรงกันข้าม นั่นคือนำ pandas DataFrame เข้าสู่ Spark cluster! SparkSession มีเมธอดสำหรับทำสิ่งนี้เช่นกัน

เมธอด .createDataFrame() รับ pandas DataFrame เป็น argument และคืนค่าเป็น Spark DataFrame

ผลลัพธ์ของเมธอดนี้จะถูกเก็บไว้ในเครื่อง ไม่ใช่ใน catalog ของ SparkSession หมายความว่าสามารถใช้เมธอดทั้งหมดของ Spark DataFrame กับมันได้ แต่จะไม่สามารถเข้าถึงข้อมูลในบริบทอื่นได้

ตัวอย่างเช่น คิวรี SQL (โดยใช้เมธอด .sql()) ที่อ้างอิงถึง DataFrame นี้จะเกิดข้อผิดพลาด หากต้องการเข้าถึงข้อมูลด้วยวิธีนี้ จะต้องบันทึกเป็น temporary table ก่อน

ทำได้โดยใช้เมธอด .createTempView() ของ Spark DataFrame ซึ่งรับ argument เพียงตัวเดียวคือชื่อของ temporary table ที่ต้องการลงทะเบียน เมธอดนี้จะลงทะเบียน DataFrame เป็นตารางใน catalog แต่เนื่องจากตารางนี้เป็นแบบ temporary จึงเข้าถึงได้เฉพาะจาก SparkSession ที่ใช้สร้าง Spark DataFrame นั้นเท่านั้น

นอกจากนี้ยังมีเมธอด .createOrReplaceTempView() ซึ่งจะสร้าง temporary table ใหม่อย่างปลอดภัยหากยังไม่มี หรืออัปเดตตารางที่มีอยู่แล้วหากมีการกำหนดไว้ก่อนหน้า เราจะใช้เมธอดนี้เพื่อหลีกเลี่ยงปัญหาจากตารางที่ซ้ำกัน

ดูแผนภาพด้านล่างเพื่อทำความเข้าใจว่าโครงสร้างข้อมูล Spark แต่ละแบบโต้ตอบกันอย่างไร

ใน workspace มี SparkSession ชื่อ spark พร้อมใช้งานแล้ว รวมถึง numpy ที่ import เป็น np และ pandas เป็น pd

แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร

พื้นฐาน PySpark

ดูคอร์ส

คำแนะนำการฝึกหัด

  • โค้ดสำหรับสร้าง pandas DataFrame ของตัวเลขสุ่มได้เตรียมไว้ให้แล้วและบันทึกอยู่ในตัวแปร pd_temp
  • สร้าง Spark DataFrame ชื่อ spark_temp โดยเรียกใช้เมธอด .createDataFrame() ของ Spark โดยส่ง pd_temp เป็น argument
  • ตรวจสอบรายการตารางใน Spark cluster และยืนยันว่า DataFrame ใหม่ ยังไม่ปรากฏ โดยใช้ spark.catalog.listTables()
  • ลงทะเบียน spark_temp DataFrame ที่เพิ่งสร้างเป็น temporary table โดยใช้เมธอด .createOrReplaceTempView() โดยตั้งชื่อ temporary table ว่า "temp" โปรดจำไว้ว่าต้องระบุชื่อตารางเป็น argument เพียงตัวเดียวของเมธอด
  • ตรวจสอบรายการตารางอีกครั้ง

แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ

ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์

# Create pd_temp
pd_temp = pd.DataFrame(np.random.random(10))

# Create spark_temp from pd_temp
spark_temp = ____

# Examine the tables in the catalog
print(____)

# Add spark_temp to the catalog
spark_temp.____

# Examine the tables in the catalog again
print(____)
แก้ไขและรันโค้ด