ข้อมูลเบื้องต้นเกี่ยวกับ Spark

แนวคิดหลักที่สนับสนุนการปรับขนาด การเพิ่มประสิทธิภาพ และการแก้ไขปัญหา อ่านสิ่งนี้ก่อนหากคุณเพิ่งเริ่มใช้ Spark in Fabric

สิ่งที่ควรทําและไม่ควรทําทั่วไป

สถานการณ์สมมติ: คุณเพิ่งเริ่มใช้ Spark สิ่งที่ควรทําและไม่ควรทําคืออะไร
กรณีการใช้งาน แนวทางปฏิบัติที่ดีที่สุด
ใช้รูปแบบอนุกรมที่ปรับให้เหมาะสม สิ่งที่ควรทํา: ต้องการรูปแบบเช่น Avro, Parquet หรือ Optimized Row Columnar (ORC) เนื่องจากมีสคีมาฝังตัว มีขนาดกะทัดรัด และเพิ่มประสิทธิภาพการจัดเก็บและการประมวลผล ใน Fabric ให้ใช้รูปแบบเดลต้าเพื่อรับประกันความเป็นอะตอม ความสม่ําเสมอ การแยก ความทนทาน (ACID) และประโยชน์ด้านประสิทธิภาพ
ระมัดระวังด้วย XML/JSON อย่าพึ่งพาการอนุมานสคีมาสําหรับไฟล์ JavaScript Object Notation (JSON) หรือ Extensible Markup Language (XML) ขนาดใหญ่ เนื่องจาก Spark จะอ่านชุดข้อมูลทั้งหมดเพื่ออนุมานสคีมา ซึ่งทําให้การประมวลผลช้าลงและใช้หน่วยความจําอย่างหนัก

ระบุ Schema หลักแบบคงที่เมื่ออ่าน JSON/XML หรือใช้ .option("samplingRatio", 0.1) เพื่อเพิ่มความเร็วในการอ่าน แต่โปรดทราบว่าหากตัวอย่างไม่ได้แสดงถึงชุดข้อมูลทั้งหมด วิธีการที่ปลอดภัยกว่าจะอนุมานสคีมาจากตัวอย่างที่เป็นตัวแทนและคงไว้สําหรับการอ่านทั้งหมด

หลีกเลี่ยงการแยกวิเคราะห์ไฟล์ XML ขนาดใหญ่ การแยกวิเคราะห์ XML ทํางานช้าลงโดยเนื้อแท้เนื่องจากการประมวลผลแท็กและการแคสต์ประเภท
เพิ่มประสิทธิภาพการรวมและการกรอง สิ่งที่ควรทํา: ใช้การตัดแต่งคอลัมน์และการกรองระดับแถวก่อนการรวมเพื่อลดการสุ่มและการใช้หน่วยความจํา

เครื่องมือเพิ่มประสิทธิภาพ Catalyst จะจัดการการกดลงเพรดิเคตโดยอัตโนมัติเมื่อคุณใช้ DataFrame API หลีกเลี่ยง Resilient Distributed Dataset (RDD) API เนื่องจากข้ามการเพิ่มประสิทธิภาพ Catalyst
ชอบ DataFrames มากกว่า RDD สิ่งที่ควรทํา: ใช้ DataFrames แทน RDD สําหรับการดําเนินการส่วนใหญ่ DataFrames ใช้เครื่องมือเพิ่มประสิทธิภาพ Catalyst และเอ็นจิ้นการดําเนินการทังสเตนเพื่อการดําเนินการที่มีประสิทธิภาพ
เปิดใช้งานการดําเนินการสืบค้นแบบปรับเปลี่ยนได้ (AQE) สิ่งที่ควรทํา: เปิด AQE เพื่อปรับพาร์ติชันแบบสุ่มให้เหมาะสมแบบไดนามิกและจัดการข้อมูลที่เบ้โดยอัตโนมัติ

การจัดการหน่วยความจํา Executor

สถานการณ์สมมติ: คุณต้องการทําความเข้าใจการจัดการหน่วยความจําตัวดําเนินการสําหรับการปรับแต่งประสิทธิภาพ

แม้ว่าตัวดําเนินการจะถูกกําหนดค่าด้วยหน่วยความจํา 56 GB แต่ Spark ก็ไม่อนุญาตให้ใช้ทั้งหมดโดยตรงสําหรับข้อมูลผู้ใช้ Spark Core แบ่งและจัดการหน่วยความจําของผู้ดําเนินการ:

  • หน่วยความจําที่สงวนไว้: ส่วนคงที่ที่สงวนไว้สําหรับค่าใช้จ่ายภายในของระบบและ Spark (เช่น Java Virtual Machine (JVM) ภายใน)

  • หน่วยความจําผู้ใช้: จัดเก็บฟังก์ชันที่ผู้ใช้กําหนด (UDF) ตัวแปรภายในเครื่อง โครงสร้างข้อมูล (รายการ แผนที่ พจนานุกรม) และวัตถุที่สร้างขึ้นระหว่างการคํานวณ

  • หน่วยความจํา: เก็บข้อมูลที่แคช/คงอยู่ ตัวแปรออกอากาศ และสับเปลี่ยนข้อมูลที่สามารถแคชได้

  • หน่วยความจําการดําเนินการ: ใช้สําหรับการคํานวณระดับกลาง (สับเปลี่ยน รวม เรียงลําดับ การรวม)

  • การแชร์หน่วยความจําแบบไดนามิก: ขอบเขตระหว่างหน่วยความจําที่เก็บข้อมูลและการดําเนินการสามารถเคลื่อนย้ายได้ Spark สามารถยืมหน่วยความจําจากภูมิภาคหนึ่งไปยังอีกภูมิภาคหนึ่ง ทําให้สามารถใช้หน่วยความจําได้อย่างยืดหยุ่น

  • หก: เกิดขึ้นเมื่อความต้องการหน่วยความจํา Storage หรือ Execution เกินหน่วยความจําที่พร้อมใช้งานหลังจากการยืม สิ่งนี้บังคับให้ข้อมูลไปยังดิสก์ ซึ่งอาจส่งผลต่อประสิทธิภาพการทํางาน

    ไดอะแกรมของการจัดการหน่วยความจํา Spark และการรั่วไหล

ข้อผิดพลาดหน่วยความจําไม่เพียงพอ (OOM)

สถานการณ์สมมติ: งาน Spark ล้มเหลวด้วยข้อผิดพลาดหน่วยความจําไม่เพียงพอ (OOM)

ไดรเวอร์ OOM:

ข้อผิดพลาด OOM ของไดรเวอร์เกิดขึ้นเมื่อไดรเวอร์ Spark เกินหน่วยความจําที่จัดสรร

สาเหตุที่พบบ่อย: การทํางานที่ต้องใช้ไดรเวอร์จํานวนมาก เช่น collect(), countByKey()หรือการ toPandas() เรียกขนาดใหญ่ที่ดึงข้อมูลมากเกินไปในหน่วยความจําของไดรเวอร์

การบรรเทาผลกระทบ: หลีกเลี่ยงการทํางานหนักของผู้ขับขี่ทุกครั้งที่ทําได้ หากหลีกเลี่ยงไม่ได้ ให้เพิ่มขนาดไดรเวอร์และเกณฑ์มาตรฐานเพื่อค้นหาการกําหนดค่าที่เหมาะสมที่สุด

ตัวดําเนินการไม่อยู่ในหน่วยความจํา (OOM):

ข้อผิดพลาด OOM ของตัวดําเนินการเกิดขึ้นเมื่อตัวดําเนินการ Spark เกินหน่วยความจําที่จัดสรร

สาเหตุที่พบบ่อย: การแปลงที่เน้นหน่วยความจําและการประมวลผลบนชุดข้อมูลขนาดใหญ่ (ตัวอย่างเช่น การรวมแบบกว้าง การรวม การสับเปลี่ยน) หรือชุดข้อมูลที่แคช/คงอยู่ซึ่งเกินหน่วยความจําที่มีอยู่ของตัวดําเนินการ (การดําเนินการ + ภูมิภาคที่เก็บข้อมูล)

การบรรเทาผลกระทบ: เพิ่มหน่วยความจําตัวดําเนินการหากจําเป็น ปรับแต่งเศษส่วนของหน่วยความจํา Spark (spark.memory.fraction, spark.memory.storageFraction) และคงอยู่อย่างเลือกสรร ตรวจสอบให้แน่ใจว่าข้อมูลที่แคชไว้พอดีกับหน่วยความจําที่พร้อมใช้งาน

ความเบ้ของข้อมูล

อาการเอียง:

  • งานบางงานใช้เวลานานกว่างานอื่นๆ ใน Spark UI (งานขั้นตอนแสดงหางหนัก)
  • ช่องว่างขนาดใหญ่ระหว่างเวลางานมัธยฐานและสูงสุดในเมตริกขั้นตอน
  • ขั้นตอนที่มีขนาดการอ่านหรือเขียนแบบสุ่มขนาดใหญ่สําหรับพาร์ติชันสองสามพาร์ติชัน

สาเหตุทั่วไป:

  • การกระจายข้อมูลที่ไม่สม่ําเสมอสําหรับคีย์รวม/กลุ่ม (ปุ่มลัด)
  • การแบ่งพาร์ติชันไม่ถูกต้องหรือพาร์ติชันน้อยเกินไปสําหรับไดรฟ์ข้อมูล
  • ความผิดปกติของข้อมูลอัปสตรีมที่สร้างเรกคอร์ดขนาดใหญ่หรือคีย์ null/empty จํานวนมาก

การบรรเทา:

  • แบ่งพาร์ติชันใหม่หรือรวมเข้าด้วยกันเพื่อเพิ่มความขนานของพาร์ติชันและขนาดสมดุล
  • ใช้การใส่เกลือคีย์หรือการแบ่งพาร์ติชันแบบกําหนดเองเพื่อกระจายปุ่มลัดข้ามพาร์ติชัน
  • ใช้ AQE (Adaptive Query Execution) เพื่อรวมพาร์ติชันหลังการสับเปลี่ยนและเปิดใช้งานการเพิ่มประสิทธิภาพการรวมแบบเอียง
  • ใช้การรวมการออกอากาศสําหรับตารางการค้นหาขนาดเล็กเพื่อหลีกเลี่ยงการสับเปลี่ยนทั้งหมด
  • คงชุดข้อมูลระดับกลางที่สมดุลไว้ก่อนขั้นตอนที่มีราคาแพงและเรียกใช้งานอีกครั้ง

แนวทางปฏิบัติที่ดีที่สุดของ UDF

สถานการณ์สมมติ: คุณต้องใช้ตรรกะแบบกําหนดเองที่ไม่สามารถแสดงผ่านฟังก์ชัน DataFrame ที่มีอยู่แล้วภายใน

ใช้ Spark DataFrame API ทุกครั้งที่ทําได้ เครื่องมือเพิ่มประสิทธิภาพ Catalyst จะปรับฟังก์ชันในตัวให้เหมาะสมและเรียกใช้บน JVM เพื่อให้มีประสิทธิภาพที่ดีที่สุด

หากคุณต้องใช้ UDF (ฟังก์ชันที่กําหนดโดยผู้ใช้) ให้หลีกเลี่ยง PySpark Python UDF ปกติ ให้พิจารณาทางเลือกต่อไปนี้แทน:

  • Pandas UDF (หรือที่เรียกว่า Vectorized UDFs): ใช้ Apache Arrow เพื่อการถ่ายโอนข้อมูลระหว่าง JVM และ Python อย่างมีประสิทธิภาพ Pandas UDF ช่วยให้สามารถทํางานแบบเวกเตอร์ ซึ่งปรับปรุงประสิทธิภาพอย่างมากเมื่อเทียบกับ Python UDF แบบแถวต่อแถว

  • Scala/Java UDFs: เรียกใช้โดยตรงบน JVM โดยหลีกเลี่ยงค่าใช้จ่ายในการทําให้เป็นอนุกรมของ Python โดยทั่วไปแล้ว Scala/Java UDF จะมีประสิทธิภาพดีกว่า Python UDF

ระมัดระวังด้วย Python UDF ผู้ดําเนินการแต่ละคนจะเปิดกระบวนการ Python แยกต่างหาก ซึ่งต้องมีการทําให้เป็นอนุกรมและการแยกข้อมูลระหว่าง JVM และ Python สิ่งนี้สร้างปัญหาคอขวดด้านประสิทธิภาพ โดยเฉพาะอย่างยิ่งในวงกว้าง 

การบันทึกข้อผิดพลาด

สถานการณ์สมมติ: แนวทางปฏิบัติที่ดีที่สุดสําหรับการบันทึกข้อผิดพลาดใน Fabric Spark
  1. ใช้ log4j แทน print() ภาระของผู้ขับขี่อย่างหนัก ด้วย log4jคุณสามารถเข้าถึงบันทึกในบันทึกไดรเวอร์และค้นหาได้ (โดยใช้ชื่อตัวบันทึก เช่น PySparkLogger)

    ไดอะแกรมของบันทึก Spark

  2. ตัดการอ่าน การเขียน และการแปลงในบล็อก try และ except ใช้สําหรับ logger.error ข้อยกเว้นและ logger.info ข้อความความคืบหน้า

    • การบันทึก Python: เหมาะอย่างยิ่งสําหรับการดําเนินการบันทึก การอัปเดตสถานะ หรือการดีบักข้อมูลจากโค้ดที่ทํางานเฉพาะบน Spark Driver โมดูลการบันทึกของ Python จะไม่เผยแพร่ไปยังบันทึกของตัวดําเนินการ ดูเอกสารประกอบการพัฒนา ดําเนินการ และจัดการสมุดบันทึก

    • สปาร์ค log4j: มาตรฐานสําหรับการบันทึกแอปพลิเคชันระดับการผลิตที่มีประสิทธิภาพใน Spark เนื่องจากผสานรวมกับบันทึกไดรเวอร์/ตัวดําเนินการของ Spark โดยกําเนิด

    ตัวอย่างการใช้งาน log4j ใน PySpark:

    import traceback
    # Get log4j logger
    log4jLogger = spark._jvm.org.apache.log4j
    logger = log4jLogger.LogManager.getLogger("PySparkLogger")
    logger.info("Application started.")
    try:
        # Create DataFrame with 20 records
        data = [(f"Name{i}", i) for i in range(1, 21)]  # 20 records
        df = spark.createDataFrame(data, ["name", "age"])
        logger.info("DataFrame created successfully with 20 records.")
        df.show(s)  # 's' is not defined -> will throw error but the application will not fail
    except Exception as e:
        logger.error(f"Error while creating or showing DataFrame: {str(e)}\n{traceback.format_exc()}")
    
  3. รวมศูนย์การตรวจสอบข้อผิดพลาด:

    • ใช้ส่วนขยายตัวส่งสัญญาณการวินิจฉัย (ตรวจสอบแอปพลิเคชัน Apache Spark ด้วย Azure Log Analytics) ในสภาพแวดล้อม และแนบกับสมุดบันทึกที่เรียกใช้แอปพลิเคชัน Spark ตัวส่งสัญญาณสามารถส่งบันทึกเหตุการณ์ บันทึกแบบกําหนดเอง (เช่น log4j) และเมตริกไปยัง Azure Log Analytics/Azure Storage/Azure Event Hubs ส่งชื่อ log4j ไปยังคุณสมบัติ: spark.synapse.diagnostic.emitter.\<destination\>.filter.loggerName.match.

    • นอกจากนี้ สําหรับการดีบัก คุณยังสามารถรวบรวมแถว/เรกคอร์ดที่ล้มเหลวไปยังตาราง Lakehouse (LH) สําหรับการเก็บข้อมูลที่ไม่ถูกต้องในระดับเรกคอร์ด