Verwirrt über die BaseSensorOperator-Parameter von Airflow: Timeout, poke_interval und mode

Sep 07 2020

Ich bin etwas verwirrt über die Funktionsweise BaseSensorOperatorder Parameter: timeout& poke_interval. Betrachten Sie diese Verwendung des Sensors:

BaseSensorOperator(
  soft_fail=True,
  poke_interval = 4*60*60,  # Poke every 4 hours
  timeout = 12*60*60,  # Timeout after 12 hours
)

In der Dokumentation wird erwähnt, dass das Zeitlimit die Aufgabe nach Ablauf auf "Fehlschlagen" setzt. Aber ich verwende ein soft_fail=True, ich glaube nicht, dass es das gleiche Verhalten beibehält, weil ich festgestellt habe, dass die Aufgabe fehlgeschlagen ist, anstatt zu überspringen, nachdem ich beide Parameter soft_failund verwendet habe timeout.

Was passiert hier?

  1. Der Sensor stößt alle 4 Stunden und wartet bei jedem Stoßen auf die Dauer des Timeouts (12 Stunden).
  2. Oder stößt es alle 4 Stunden für insgesamt 3 Stöße und tritt dann eine Zeitüberschreitung auf?
  3. Was passiert mit diesen Parametern, wenn ich mode = "reschedule" verwende?

Hier ist die Dokumentation des BaseSensorOperator

class BaseSensorOperator(BaseOperator, SkipMixin):
    """
    Sensor operators are derived from this class and inherit these attributes.
    Sensor operators keep executing at a time interval and succeed when
    a criteria is met and fail if and when they time out.
    :param soft_fail: Set to true to mark the task as SKIPPED on failure
    :type soft_fail: bool
    :param poke_interval: Time in seconds that the job should wait in
        between each tries
    :type poke_interval: int
    :param timeout: Time, in seconds before the task times out and fails.
    :type timeout: int
    :param mode: How the sensor operates.
        Options are: ``{ poke | reschedule }``, default is ``poke``.
        When set to ``poke`` the sensor is taking up a worker slot for its
        whole execution time and sleeps between pokes. Use this mode if the
        expected runtime of the sensor is short or if a short poke interval
        is requried.
        When set to ``reschedule`` the sensor task frees the worker slot when
        the criteria is not yet met and it's rescheduled at a later time. Use
        this mode if the expected time until the criteria is met is. The poke
        inteval should be more than one minute to prevent too much load on
        the scheduler.
    :type mode: str
    """

Antworten

6 y2k-shubham Sep 07 2020 at 11:44

Begriffe definieren

  1. poke_interval: die Dauer s / w aufeinanderfolgender "Stöße" (Bewertung der notwendigen Bedingung, die "erfasst" wird)

  2. timeout: Nur auf unbestimmte Zeit zu stöbern ist unzulässig (wenn beispielsweise Ihr Buggy-Code am Tag stochert und 29 wird, wenn der Monat 2 ist, stochert er bis zu 4 Jahre lang). Wir definieren also eine maximale Zeitspanne, nach der wir aufhören zu stochern und beenden (der Sensor ist entweder FAILEDoder markiert SKIPPED).

  3. soft_fail: Normalerweise (wann soft_fail=False) wird der Sensor als FAILEDnach Zeitüberschreitung markiert . Wenn soft_fail=True, wird der Sensor stattdessen als SKIPPEDnach Zeitüberschreitung markiert

  4. mode: Dies ist ein wenig komplex

    • Jede Aufgabe (einschließlich Sensor), die ausgeführt wird, verschlingt eine slotin einem Pool (entweder defaultPool oder explizit angegeben pool). Dies bedeutet im Wesentlichen, dass einige Ressourcen benötigt werden.
    • Für Sensoren ist dies
      • verschwenderisch : da ein Slot auch dann verbraucht wird, wenn wir nur warten (keine eigentliche Arbeit erledigen)
      • gefährlich : Wenn Sie Ihren Workflow zu viele Sensoren , die in gehen hat Abtasten der gleichen Zeit um, sie eine Menge von Ressourcen für eine ganze Bit gefrieren kann. Tatsächlich sind zu viele mit ExternalTaskSensors dafür berüchtigt , ganze Workflows (DAGs) in Deadlocks zu stecken
    • Um dieses Problem zu lösen, hat Airflow v1.10.2 mode s in Sensoren eingeführt
      • mode='poke' (Standard) bedeutet das vorhandene Verhalten, das wir oben besprochen haben
      • mode='reschedule'Mittel nach einem Poke Versuch , anstatt zu schlafen geht , wird der Sensor so verhalten , als ob es (in der aktuellen Versuch) gescheitert und Status wird von ändern RUNNINGzu UP_FOR_RETRY. Auf diese Weise wird der Steckplatz freigegeben , sodass andere Aufgaben ausgeführt werden können, während auf einen weiteren Poke-Versuch gewartet wird
    • Hier wird das relevante Snippet aus dem Code zitiert
    if self.reschedule:
        reschedule_date = timezone.utcnow() + timedelta(
            seconds=self._get_next_poke_interval(started_at, try_number))
        raise AirflowRescheduleException(reschedule_date)
    else:
        sleep(self._get_next_poke_interval(started_at, try_number))
        try_number += 1
    
    • Weitere Informationen finden Sie im Abschnitt Sensors Params

Und jetzt beantworten Sie Ihre Fragen direkt

Q1

  1. Der Sensor stößt alle 4 Stunden und wartet bei jedem Stoßen auf die Dauer des Timeouts (12 Stunden).
  2. Oder stößt es alle 4 Stunden für insgesamt 3 Stöße und tritt dann eine Zeitüberschreitung auf?

Punkt 2. ist richtig

Q2

Was passiert mit diesen Parametern, wenn ich mode = "reschedule" verwende?

Wie bereits erläutert, ist jeder dieser Parameter unabhängig und die Einstellung mode='reschedule'ändert nichts an seinem Verhalten

1 SanajaobaThongram Oct 12 2020 at 10:31
BaseSensorOperator(
  soft_fail=True,
  poke_interval = 4*60*60,  # Poke every 4 hours
  timeout = 12*60*60,  # Timeout of 12 hours
  mode = "reschedule"
)

Angenommen, die Kriterien werden beim ersten Mal nicht erfüllt. So läuft es nach 4 Stunden wieder. Der Worker-Slot wird jedoch während des Wartens freigegeben, da wir den verwenden mode="reschedule".

Das habe ich verstanden.