كيف ننظم معالجة البيانات باستخدام Apache Airflow

مرحبا! اسمي نيكيتا فاسيليوك ، أنا مهندس بيانات في قسم البيانات والتحليلات في لامودا. في قسمنا ، يلعب Airflow دور منسق عمليات معالجة البيانات الضخمة ، بمساعدته نقوم بتحميل البيانات من الأنظمة الخارجية إلى Hadoop ، وتدريب نماذج ML ، وكذلك إجراء فحوصات جودة البيانات ، وحسابات أنظمة التوصية ، ومقاييس مختلفة ، واختبارات A / B وأكثر من ذلك بكثير. ...



صورة



سأشرح في هذا المقال:



  • ما نوع الوحش هذا تدفق الهواء ، وما المكونات التي يتكون منها وكيف يتفاعلون مع بعضهم البعض
  • حول الكيانات الرئيسية لـ Airflow: خطوط الأنابيب المسماة DAG ، المشغل وعدد قليل من الأشياء الأخرى
  • كيف تنجح في تطوير Airflow
  • كيف نفذنا توليد خطوط الأنابيب وما يسمى "الكتابة التصريحية لخطوط الأنابيب"
  • حول إيجابيات وسلبيات استخدام Airflow


ما هو تدفق الهواء



Airflow هو عبارة عن منصة لإنشاء خطوط الأنابيب ومراقبتها وتنظيمها. تم إنشاء هذا المشروع مفتوح المصدر ، المكتوب بلغة Python ، في عام 2014 في Airbnb. في عام 2016 ، دخلت Airflow تحت جناح مؤسسة Apache Software Foundation ، وخضعت لحاضنة ، وفي بداية عام 2019 أصبح مشروع Apache عالي المستوى.



في عالم معالجة البيانات ، يسميها البعض أداة ETL ، ولكن هذا ليس بالضبط ETL بالمعنى الكلاسيكي ، مثل Pentaho و Informatica PowerCenter و Talend وغيرها. إن تدفق الهواء هو منظم ، "كرون على البطاريات": فهو لا يقوم بالعمل الشاق لنقل البيانات ومعالجتها بنفسه ، ولكنه يخبر الأنظمة والأطر الأخرى بما يجب القيام به ويراقب حالة التنفيذ. نحن نستخدمه بشكل أساسي لتشغيل الاستعلامات في وظائف Hive أو Spark.



المفسد
Airflow, worker ( ), . , , .



لا يقتصر نطاق المهام التي تم حلها باستخدام Airflow على تشغيل شيء ما في مجموعة Hadoop. يمكنه تشغيل كود Python ، وتنفيذ أوامر Bash ، واستضافة حاويات Docker و pods في Kubernetes ، وتنفيذ استعلامات على قاعدة البيانات المفضلة لديك ، والمزيد.



بنية تدفق الهواء



صورة



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



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



مكونات تدفق الهواء



خادم الويب



Webserver هو واجهة ويب تعرض ما يحدث مع خط الأنابيب. هذه الصفحة



صورة



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



إذا نقرنا على خط الأنابيب ، فسوف نقع في قائمة عرض الرسم البياني. يتم عرض المهام والروابط بينهما هنا.



صورة



توجد قائمة Tree View بجوار عرض الرسم البياني. تم إنشاؤه لإعادة المهام وعرض الإحصائيات والسجلات. يتم عرض الرسم البياني الشبيه بالشجرة على الجانب الأيسر ، مقابله يوجد جدول به محفوظات بدء المهمة.



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



صورة



المجدول - كما يوحي الاسم ، يطلق خطوط الأنابيب عندما يحين وقتها. إنها عملية Python التي تنتقل بشكل دوري إلى الدليل مع خطوط الأنابيب ، وتسحب حالتها الحالية من هناك ، وتتحقق من الحالة ، وتبدأها. بشكل عام ، يعد المجدول الأكثر إثارة للاهتمام وفي نفس الوقت هو عنق الزجاجة في بنية Airflow.



  • التحذير الأول هو أنه يمكن تشغيل مثيل جدولة واحد فقط في كل مرة. هذا يعني أن وضع High Availability غير ممكن حاليًا (يخطط المطورون لإضافة Scheduler HA إلى Airflow الإصدار 2.0).
  • : , - . , - , .


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



العامل



العامل هو المكان الذي يعمل فيه كودنا ويتم إنجاز المهام. يدعم تدفق الهواء العديد من المنفذين:



  • الأول ، الأبسط ، هو SequentialExecutor. يقوم بتشغيل المهام الواردة بالتسلسل ، وإيقاف المجدول مؤقتًا لمدة تنفيذها.
  • LocalExecutor , , LocalExecutor . : - SQLite, LocalExecutor SequentialExecutor.
  • CeleryExecutor , . Celery – , RabbitMQ Redis. , .
  • DaskExecutor Dask – .
  • KubernetesExecutor pod Kubernetes.
  • DebugExecutor IDE.


كيانات Apache Airflow



خط الأنابيب ، أو DAG



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



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



dag = DAG(
   dag_id="load_some_data",
   schedule_interval="0 1 * * *",
   default_args={
       "start_date": datetime(2020, 4, 20),
       "owner": "DE",
       "depends_on_past": False,
       "sla": timedelta(minutes=45),
       "email": "<your_email_here>",
       "email_on_failure": True,
       "retries": 2,
       "retry_delay": timedelta(minutes=5)
   }
)


يحتوي dag_id على الاسم الفريد لخط الأنابيب. بعد ذلك ، نستخدم Schedule_interval لتحديد عدد مرات تشغيله.



نقطة مهمة للغاية: نظرًا لأن Airflow تم تطويره بواسطة شركة دولية ، فإنه يعمل فقط في UTC. في الوقت الحالي ، لا توجد طريقة عقلانية لجعل Airflow يعمل في منطقة زمنية مختلفة ، لذلك عليك أن تتذكر باستمرار الفرق بين منطقتنا الزمنية والتوقيت العالمي المنسق (UTC). في الإصدار 1.10.10 ، أصبح من الممكن تغيير المنطقة الزمنية في واجهة المستخدم ، ولكن هذا ينطبق فقط على واجهة الويب ، وستظل خطوط الأنابيب تعمل بالتوقيت العالمي المنسق.



المعلمة default_args هي قاموس يصف الوسيطات الافتراضية لجميع المهام ضمن مسار الأنابيب هذا. أسماء معظم المعلمات تصف نفسها جيدًا ، ولن أسهب في الحديث عنها.



المشغل أو العامل



عامل التشغيل هو فئة Python التي تصف الإجراءات التي يجب القيام بها في مهمتنا اليومية من أجل إسعاد المحلل.



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



run_sql = HiveOperator(
   dag=dag,
   task_id="run_sql",
   hive_cli_conn_id="hive",
   hql="""
       INSERT OVERWRITE TABLE some_table
       SELECT * FROM other_table t1
       JOIN another_table t2 on ...
       WHERE other_table.dt = '{{ ds }}'
   """
)

notify = SlackAPIPostOperator(
   dag=dag,
   task_id="notify_slack",
   slack_conn_id="slack",
   token=token,
   channel="airflow_alerts",
   text="Guys, I'm done for {{ ds }}"
)

run_sql >> notify


يوجد جزء من قالب Jinja في الطلب نمرره إلى مُنشئ المشغل. Jinja هي مكتبة قوالب بايثون.



يخزن كل إطلاق خط أنابيب معلومات حول تاريخ الإطلاق. تقع في متغير يسمى تاريخ التنفيذ. {{ds}} هو ماكرو سيستغرق فقط التاريخ بالتنسيق٪ Y-٪ m-٪ d في تاريخ التنفيذ. في لحظة معينة قبل بدء المشغل ، ستعرض Airflow سلسلة استعلام ، وتستبدل التاريخ المطلوب هناك وترسل طلبًا للتنفيذ.



ds ليس الماكرو الوحيد ، فهناك حوالي 20 منهم (قائمة بجميع وحدات الماكرو المتوفرة) . وهي تتضمن تنسيقات تاريخ مختلفة واثنين من الوظائف للعمل مع التواريخ - إضافة أو طرح بعض الأيام.



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



على سبيل المثال ، إذا أردنا إعادة تشغيل خط الأنابيب ليوم الثلاثاء الماضي ، فعند استخدام datetime.now () ، سنقوم بالفعل بإعادة حساب خط الأنابيب لهذا اليوم ، وليس للتاريخ المطلوب. بالإضافة إلى ذلك ، قد لا تكون بيانات اليوم جاهزة في هذه المرحلة.



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



العاطفة



من المستحيل التحدث عن Airflow دون ذكر العاطفة. فقط في حالة ما ، دعني أذكرك: إن idempotency هي خاصية لكائن ، عندما تعيد تطبيق عملية على كائن ، تُرجع دائمًا نفس النتيجة.



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



تم تطوير تدفق الهواء كأداة لحل مهام معالجة البيانات. في هذا العالم ، نقوم عادةً بمعالجة جزء كبير من البيانات فقط عندما تكون جاهزة ، أي في اليوم التالي. وقد وضع مبتكرو Airflow في الأصل مثل هذا المفهوم في منتجاتهم.



صورة



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



مشغلي تدفق الهواء الأكثر شيوعًا



في Airflow ، لا يوجد فقط المشغلون الذين يذهبون إلى Hive ويرسلون شيئًا ما للتراخي. في الواقع ، هناك الكثير من المشغلين هناك. في المقال ، أخرجت أكثرها شعبية وفائدة.



  • BashOpetator و PythonOperator. كل شيء واضح معهم: يرسلون أمر bash ووظيفة python للتنفيذ ، على التوالي.
  • هناك مجموعة كبيرة ومتنوعة من المشغلين لتقديم الاستفسارات إلى قواعد البيانات المختلفة. يتم دعم معايير Postgres و MySQL و Oracle و Hive و Presto. إذا لم يكن هناك عامل لقاعدة البيانات المفضلة لديك لسبب ما ، فيمكنك استخدام JdbcOperator أكثر عمومية أو كتابة الخاصة بك ، Airflow يسمح بذلك.
  • Sensor – , . , - . , , . , : 3 , . . , , .
  • BranchPythonOperator – , , python , , .
  • DockerOpetator Docker- . , Docker- , . , .
  • KubernetesPodOperator pod Kubernetes.
  • DummyOperator , .


Lamoda



  • – LamodaDockerOperator. , : - Hadoop, . LamodaDockerOperator Spark- , python.
  • LamodaHiveperator – , . Hive. , - , , . , , HiveCliHook HiveServer2Hook, .
  • – ExternalTaskSensor. . , Hadoop . , , , - , , . , - HDFS, Airflow.
  • BashOperator, PythonOperator – , bash- python .
  • , . - , .


Airflow



  • Variables – , , , . , . , Hive, HDFS, . dev- prod-, .
  • Connections – , . Airflow : http ftp, .
  • Hooks – , .
  • SLA -. , . SLA , , - - . - : - , Airflow .
  • – XCom, cross-communication. : , json-. – 48 .
  • – , . , . , 5, , , , .


صورة



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



صورة



كيف تنجح في تطوير تدفق الهواء



فيما يلي بعض النصائح لمساعدتك على تجنب التعرض لطلقة في القدم عند استخدام Airflow:



  • من المفيد الاحتفاظ بكل خط أنابيب (أو مولد خط أنابيب ، أكثر من ذلك أدناه) في ملف منفصل. أعرف على الفور الملف الذي أحتاج إلى الانتقال إليه لإلقاء نظرة على خط الأنابيب أو المولد المطلوب.
  • , , . , -, . , - , . : , , .
  • – schedule_interval start_date dag_id. , Airflow , - -. DAGS , Scheduler, . , , dag_id. , .
  • catchup. True, Airflow , start_date . , . False Airflow . , Airflow True ( -).
  • – . , python , airflow DAG, , DAG. . , , . REST API, requests.get() .


:



منذ بداية استخدام Airflow ، أبقينا تكوينات خطوط الأنابيب منفصلة عن الكود. في البداية ، كان هذا بسبب خصوصيات خطة الانتشار ، ولكن تدريجيًا ترسخ هذا النهج. والآن نستخدم التكوينات أينما كان هناك تلميح من المتداول. بالنسبة لنا ، يتعلق هذا بوظائف Spark التي نديرها من Docker. من هنا جاءت القصة مع الكتابة التوضيحية لخطوط الأنابيب.



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



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



from libs.dag_from_config.dag_generator import DagGenerator
from libs.runners.docker_runner import DockerRunner

generator = DagGenerator(config_dir='dag_configs/docker_runner', prefix='docker')
dags = generator.generate(task_runner=DockerRunner)

for dag in dags:
   globals()[dag.dag_id] = dag  #     


هذا ما يبدو عليه ملف التكوين النموذجي. لوصف التكوينات ، نستخدم تنسيق HOCON ، وهو مجموعة شاملة من JSON. وهو يدعم واردات ملفات HOCON الأخرى ويمكن أن يشير إلى قيم المتغيرات الأخرى.



في التكوين على مستوى خط الأنابيب (كتلة الإحالة) ، يمكنك تحديد العديد من المعلمات ، ولكن أهمها الاسم وتاريخ البدء والجدول الزمني.



docker_image = "docker_registry/attribution/calculation:1.1.0"

dags {
 attribution {
   owner = "RND"
   name = "attribution"
   start_date = "20190601"
   emails = [...]
   schedule_interval = "0 1 * * *"
   depends_on_past = true
   concurrency = 4

   description = """
   -    z_log
   -        
   -  ,    
   -     
   """

   tags = ["critical"]


هنا يمكنك تحديد التزامن - كم عدد المهام التي سيتم تشغيلها في وقت واحد في تشغيل واحد. لقد أضفنا مؤخرًا كتلة تحتوي على وصف مختصر لخط الأنابيب هنا. بعد ذلك ، سيتم إرساله ، مع بقية المعلومات حول خط الأنابيب ، إلى Confluence (قمنا بتنفيذ الإرسال باستخدام Foliant ). اتضح أنها مريحة للغاية: بهذه الطريقة نوفر الوقت للمطورين المحفورين لإنشاء صفحات في منطقة التقاء.



بعد ذلك يأتي الجزء المسؤول عن تشكيل المهام. أولاً ، في كتلة الاتصالات ، نشير إلى أي اتصال في Airflow نحتاج إلى أخذ معلمات للاتصال بمصدر خارجي - في المثال ، هذا هو DWH الخاص بنا.



docker {
 connections {
   LMD_DWH = "dwh"
 }

 containers {
   desktop {
     image = ${docker_image}
     connections = [LMD_DWH]

     environment {
       LMD_YARN_QUEUE = "{{ var.value.YARN_QUEUE }}"
       LMD_INSTANCES = 60
       LMD_MEMORY_PER_INSTANCE = "4g"
       LMD_ZLOG_SOURCE = "z_log_db.z_log"
       LMD_ATTRIBUTION_TABLE = "{{ var.value.HIVE_DB_DERIVATIVES }}.z_log_attribution"
       LMD_ORDERS_TABLE = "rocket_dwh_bl.fct_orderitem_detail"
       LMD_PLATFORMS = "desktop"

       LMD_RUN_DATE = "{{ ds_nodash }}"
     }
   }
   mobile {...}
   iOS {...}
   Android {...}
 }
 tasks = [[desktop, mobile, iOS, Android]]
}


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



قد تلاحظ أن قوالب Jinja تظهر في قيم بعض متغيرات البيئة. لتحديد قائمة انتظار في YARN ، نستخدم صيغة Airflow القياسية لاسترداد القيم المتغيرة. للإشارة إلى تاريخ الإطلاق ، نستخدم الماكرو {{ds_nodash}} ، والذي يمثل تاريخ تاريخ التنفيذ بدون واصلات. يحتوي التكوين على 3 مهام أخرى مماثلة ، يتم إخفاؤها من أجل الوضوح.



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



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



  awaits {
     z_log_compaction {
       dag = "compactor_daily_23_21_A_A_A"
       task = "compact_z_log_db_z_log"
       timedelta = 3hr37m
     }
     oracle_bl_fct_orderitem_detail {
       dag = "await_bl_fct_orderitem_detail_0_1_A_A_A"
     }
   }
 }
}


النص الكامل لملف التكوين
docker_image = "docker_registry/attribution/calculation:1.1.0"

dags {
 attribution {
   owner = "RND"
   name = "attribution"
   start_date = "20190601"
   emails = [...]
   schedule_interval = "0 1 * * *"
   depends_on_past = true
   concurrency = 4

   description = """
   -    z_log
   -        
   -  ,    
   -     
   """

   tags = ["critical"]


   docker {
     connections {
       LMD_DWH = "dwh"
     }

     containers {
       desktop {
         image = ${docker_image}
         connections = [LMD_DWH]

         environment {
           LMD_YARN_QUEUE = "{{ var.value.YARN_QUEUE }}"
           LMD_INSTANCES = 60
           LMD_MEMORY_PER_INSTANCE = "4g"
           LMD_ZLOG_SOURCE = "z_log_db.z_log"
           LMD_ATTRIBUTION_TABLE = "{{ var.value.HIVE_DB_DERIVATIVES }}.z_log_attribution"
           LMD_ORDERS_TABLE = "rocket_dwh_bl.fct_orderitem_detail"
           LMD_PLATFORMS = "desktop"

           LMD_RUN_DATE = "{{ ds_nodash }}"
         }
       }
       mobile {...}
       iOS {...}
       Android {...}
     }
     tasks = [[desktop, mobile, iOS, Android]]
   }


   awaits {
     z_log_compaction {
       dag = "compactor_daily_23_21_A_A_A"
       task = "compact_z_log_db_z_log"
       timedelta = 3hr37m
     }
     oracle_bl_fct_orderitem_detail {
       dag = "await_bl_fct_orderitem_detail_0_1_A_A_A"
     }
   }
 }
}




هذا ما نحصل عليه بعد جيل:



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


صورة



ماذا نريد أن نفعل بعد ذلك



أولاً ، قم ببناء بيئة ميزة كاملة. لدينا الآن منصة تطوير واحدة لاختبار جميع خطوط الأنابيب لدينا. وقبل الاختبار ، تحتاج إلى التأكد من أن منظر مطور البرامج مجاني الآن.



في الآونة الأخيرة ، توسع فريقنا ، وزاد عدد المتقدمين. لقد وجدنا حلاً مؤقتًا للمشكلة وأخبرنا الآن في Slack عندما نتعامل مع dev. إنه يعمل ، لكنه لا يزال يمثل عنق الزجاجة في التطوير والاختبار.



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



الخيار الثاني لتنفيذ بيئة الميزة هو تنظيم مستودع بفرع تطوير مشترك ، حيث يتم دمج كود المطورين ونشره تلقائيًا في بيئة التطوير. الآن نحن نتطلع بنشاط نحو هذا المخطط.



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



لماذا اخترنا Airflow



  1. أولاً ، هذا هو Python ، حيث بمساعدة حلقتين واثنين من الشروط ، يمكنك إنشاء خط أنابيب أنيق يعمل بشكل صحيح. ولن تحتاج إلى وصفها في جزء كبير من XML. بالإضافة إلى ذلك ، يتوفر نظام Python البيئي بالكامل تقريبًا وحديقة حيوانات المكتبات بأكملها خارج الصندوق ، والتي يمكن استخدامها كما تريد.
  2. يبسط غياب XML مراجعة التعليمات البرمجية بشكل كبير. كتبنا رمز خط الأنابيب والتكوينات له ، وكل شيء رائع ، كل شيء يعمل. في الواقع ، يمكنك السحب بتنسيق XML أو أي تنسيق تكوين آخر ، ولكن هذه بالفعل مسألة ذوق.
  3. unit-, , .
  4. , «», . Airflow . , , .
  5. Airflow ( ).
  6. Active Directory RBAC (role-based access control, )
  7. Worker Celery Kubernetes.
  8. open source-, , .
  9. Airflow , . .
  10. : statsd , Sentry – , Airflow , . Airflow-exporter Prometheus.


Airflow,



  1. – : , , execution_date – , .
  2. - -, , , Apache NiFi. – code-review diff- , .
  3. Scheduler - .
  4. – , . – .
  5. Airflow : . , , . RBAC ( ) , UI (, , ). RBAC – security Flask, .
  6. : , , -, , . , .


Airflow



  • crontab’a cron .
  • Python.
  • - Docker, , .
  • , , real time.
  • Airflow , “, , , Z – ”.


Airflow



  • Astronomer, hosted- Airflow Kubernetes. –
  • Astronomer Airflow –
  • Airflow () Slack ().



All Articles