نظرة عامة على واجهة المستخدم الجديدة للبث المنظم في Apache Spark ™ 3.0

تم إعداد ترجمة المقال عشية بدء دورة مهندس البيانات .










تم تقديم البث المنظم لأول مرة في Apache Spark 2.0. أثبتت هذه المنصة نفسها كأفضل خيار لبناء تطبيقات البث الموزعة. يعمل توحيد واجهة برمجة تطبيقات SQL / Dataset / DataFrame ووظائف Spark المدمجة على تسهيل تنفيذ أساسياتهم المعقدة مثل تجميع الدفق والانضمام إلى البث ودعم النوافذ. منذ إصدار البث المهيكل ، كان طلبًا شائعًا من المطورين لتحسين التحكم في البث ، تمامًا كما فعلنا في Spark Streaming (مثل DStream). في Apache Spark 3.0 ، أصدرنا واجهة مستخدم جديدة للبث المنظم.



يوفر البث الهيكلي لواجهة المستخدم طريقة سهلة لمراقبة جميع مهام البث من خلال رؤى وإحصاءات قابلة للتنفيذ ، مما يسهل استكشاف المشكلات وإصلاحها أثناء تصحيح الأخطاء وتحسين رؤية الإنتاج باستخدام مقاييس الوقت الفعلي. تقدم واجهة المستخدم مجموعتين من الإحصائيات: 1) معلومات مجمعة حول مهمة الاستعلام المتدفق و 2) معلومات إحصائية مفصلة حول طلبات التدفق ، بما في ذلك معدل الإدخال ، ومعدل العملية ، وصفوف الإدخال ، ومدة الدفعة ، ومدة العملية ، إلخ.



معلومات مجمعة حول تدفق مهام الاستعلام



عندما يرسل مطور استعلام SQL متدفقًا ، فإنه يظهر في علامة التبويب الهيكلية ، والتي تتضمن كلاً من استعلامات البث النشطة والاستعلامات المكتملة. سيوفر جدول النتائج بعض المعلومات الأساسية المتعلقة بطلبات البث ، بما في ذلك اسم الطلب والحالة والمعرف ومعرف التشغيل ووقت الإرسال ومدة الطلب ومعرف الحزمة الأخير ، بالإضافة إلى المعلومات المجمعة مثل متوسط ​​معدل الاستلام ومتوسط ​​معدل المعالجة. هناك ثلاثة أنواع من حالة طلب البث: قيد التشغيل ، وتم الانتهاء ، وفشل. يتم سرد جميع الطلبات المنتهية والفاشلة في جدول طلبات البث المكتمل. يعرض عمود الخطأ تفاصيل استثناء الطلب الفاشل.







يمكننا عرض الإحصائيات التفصيلية لطلب التدفق من خلال النقر على رابط معرف التشغيل.



معلومات إحصائية مفصلة



تعرض صفحة الإحصائيات المقاييس بما في ذلك معدل الاستيعاب / المعالجة ووقت الاستجابة ومدة التشغيل التفصيلية ، والتي تعد مفيدة لفهم حالة طلبات البث ، مما يسهل تصحيح الأخطاء في معالجة الطلبات.









يحتوي على المقاييس التالية:



  • معدل الإدخال : المعدل الإجمالي (عبر جميع المصادر) لوصول البيانات.
  • معدل العملية : المعدل المجمع (عبر جميع المصادر) الذي يعالج به Spark البيانات.
  • مدة الدفعة : مدة كل دفعة.
  • مدة العملية : الوقت المستغرق لإجراء عمليات مختلفة بالمللي ثانية.


المعاملات المرصودة مذكورة أدناه:



  • addBatch: الوقت المستغرق في قراءة بيانات الإدخال للدفعة الصغيرة من المصادر ومعالجتها وكتابة بيانات إخراج الدُفعة للمزامنة. يستغرق هذا عادةً معظم وقت الدفعة الجزئية.
  • getBatch: الوقت المستغرق لإعداد طلب منطقي لقراءة بيانات إدخال العبوة الدقيقة الحالية من المصادر.
  • getOffset: الوقت المستغرق في الاستعلام عن المصادر إذا كانت لديها مدخلات جديدة.
  • walCommit: يكتب الإزاحة في سجلات البيانات الوصفية.
  • queryPlanning: إنشاء خطة تنفيذ.


وتجدر الإشارة إلى أنه لن يتم عرض جميع العمليات المدرجة في واجهة المستخدم. هناك عمليات مختلفة مع أنواع مختلفة من مصادر البيانات ، لذلك يمكن إجراء بعض العمليات المدرجة في طلب دفق واحد.



استكشاف أخطاء تدفق الأداء باستخدام واجهة المستخدم وإصلاحها



في هذا القسم ، سنلقي نظرة على بعض الحالات التي يشير فيها البث المنظم لواجهة المستخدم الجديدة إلى حدوث شيء خارج عن المألوف. يبدو طلب العرض التوضيحي عالي المستوى هكذا ، وفي كل حالة سنفترض بعض الشروط المسبقة:



import java.util.UUID

val bootstrapServers = ...
val topics = ...
val checkpointLocation = "/tmp/temporary-" + UUID.randomUUID.toString

val lines = spark
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", bootstrapServers)
    .option("subscribe", topics)
    .load()
    .selectExpr("CAST(value AS STRING)")
    .as[String]

val wordCounts = lines.flatMap(_.split(" ")).groupBy("value").count()

val query = wordCounts.writeStream
    .outputMode("complete")
    .format("console")
    .option("checkpointLocation", checkpointLocation)
    .start()


زيادة الكمون بسبب قوة المعالجة غير الكافية



في الحالة الأولى ، نقوم بتشغيل طلب لمعالجة بيانات Apache Kafka في أسرع وقت ممكن. لكل دفعة ، تعالج وظيفة البث جميع البيانات المتاحة في كافكا. إذا كانت قوة المعالجة غير كافية للتعامل مع بيانات الاندفاع ، فسوف يزداد زمن الانتقال بسرعة. الحكم الأكثر بديهية هو أن صفوف الإدخال ومدة الدُفعات ستنمو خطيًا. تحدد معلمة "صفوف الإدخال" أن وظيفة الدفق يمكنها معالجة 8000 عملية كتابة بحد أقصى في الثانية. لكن معدل الإدخال الحالي يبلغ حوالي 20000 سجل في الثانية. يمكننا تزويد وظيفة الخيوط بمزيد من الموارد للتشغيل ، أو يمكننا إضافة أقسام كافية للتعامل مع جميع المستهلكين اللازمين لمواكبة المنتجين.







كمون مستقر ولكن مرتفع



كيف تختلف هذه الحالة عن سابقتها؟ لا يزيد زمن الانتقال ، لكنه يظل مستقرًا ، كما هو موضح في لقطة الشاشة التالية:







وجدنا أن معدل العملية يمكن أن يظل ثابتًا عند نفس معدل الإدخال. هذا يعني أن قوة معالجة الوظيفة كافية لمعالجة بيانات الإدخال. ومع ذلك ، فإن وقت المعالجة لكل دفعة ، أي التأخير ، لا يزال 20 ثانية. السبب الرئيسي لارتفاع وقت الاستجابة هو وجود الكثير من البيانات في كل دفعة. يمكننا عادة تقليل زمن الوصول عن طريق زيادة التوازي لهذه الوظيفة. بعد إضافة 10 أقسام أخرى من كافكا و 10 نوى لمهام Spark ، وجدنا أن زمن الانتقال يبلغ حوالي 5 ثوانٍ - أفضل بكثير من 20 ثانية.







استخدم مخطط مدة العملية لاستكشاف الأخطاء وإصلاحها



يعرض مخطط مدة العملية مقدار الوقت المستغرق في تنفيذ عمليات مختلفة بالمللي ثانية. هذا مفيد لفهم توقيت كل دفعة وتسهيل استكشاف الأخطاء وإصلاحها. دعنا نستخدم عمل تحسين الأداء " SPARK-30915 : تجنب قراءة ملف سجل البيانات الوصفية عند البحث عن أحدث معرّف دفعة" في مجتمع Apache Spark كمثال.

قبل هذا التحسين ، كانت كل دفعة لاحقة بعد الضغط تستغرق وقتًا أطول من الدُفعات الأخرى ، عندما يصبح سجل البيانات الوصفية المضغوط ضخمًا.







بعد فحص الكود ، تم العثور على قراءة غير ضرورية لملف السجل المضغوط وتم إصلاحها. يؤكد الرسم البياني التالي مدة العملية التأثير المتوقع:







خطط للمستقبل



كما هو موضح أعلاه ، سيساعد البث المنظم لواجهة المستخدم الجديدة المطورين على التحكم بشكل أفضل في وظائف البث من خلال الحصول على معلومات أكثر فائدة حول طلبات البث. كإصدار مبكر ، لا تزال واجهة المستخدم الجديدة قيد التطوير وسيتم تحسينها في الإصدارات المستقبلية. هناك العديد من الميزات التي يمكن تنفيذها في المستقبل غير البعيد ، بما في ذلك على سبيل المثال لا الحصر ما يلي:



  • تعرف على المزيد حول تنفيذ طلب التدفق: البيانات المتأخرة والعلامات المائية ومقاييس حالة البيانات والمزيد.
  • دعم واجهة المستخدم المتدفقة الهيكلية على خادم Spark History.
  • المزيد من الدلائل الملحوظة للسلوك غير المعتاد: الكمون ، إلخ.


جرب واجهة مستخدم جديدة



جرب واجهة Spark Streaming UI الجديدة في Apache Spark 3.0 في Databricks Runtime 7.1 الجديد. إذا كنت تستخدم دفاتر Databricks ، فسيوفر لك هذا أيضًا طريقة سهلة لمراقبة حالة أي طلب دفق في دفتر الملاحظات وإدارة طلباتك . يمكنك التسجيل للحصول على حساب مجاني في Databricks والبدء في دقائق مجانًا ، دون أي معلومات ائتمانية.






جودة البيانات في DWH هي اتساق مستودع البيانات. ندوة مجانية على الويب.






اقتراحات للقراءة:



أداة بناء البيانات ، أو ما يشترك فيه مستودع البيانات و Smoothie

في الغوص في Delta Lake: Schema Enforcement and Evolution

High Speed ​​Apache Parquet in Python with Apache Arrow



All Articles