Verwirrt über die BaseSensorOperator-Parameter von Airflow: Timeout, poke_interval und mode
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?
- Der Sensor stößt alle 4 Stunden und wartet bei jedem Stoßen auf die Dauer des Timeouts (12 Stunden).
- Oder stößt es alle 4 Stunden für insgesamt 3 Stöße und tritt dann eine Zeitüberschreitung auf?
- 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
Begriffe definieren
poke_interval: die Dauer s / w aufeinanderfolgender "Stöße" (Bewertung der notwendigen Bedingung, die "erfasst" wird)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 entwederFAILEDoder markiertSKIPPED).soft_fail: Normalerweise (wannsoft_fail=False) wird der Sensor alsFAILEDnach Zeitüberschreitung markiert . Wennsoft_fail=True, wird der Sensor stattdessen alsSKIPPEDnach Zeitüberschreitung markiertmode: Dies ist ein wenig komplex- Jede Aufgabe (einschließlich Sensor), die ausgeführt wird, verschlingt eine
slotin einem Pool (entwederdefaultPool oder explizit angegebenpool). 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
modes in Sensoren eingeführtmode='poke'(Standard) bedeutet das vorhandene Verhalten, das wir oben besprochen habenmode='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 ändernRUNNINGzuUP_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
- Jede Aufgabe (einschließlich Sensor), die ausgeführt wird, verschlingt eine
Und jetzt beantworten Sie Ihre Fragen direkt
Q1
- Der Sensor stößt alle 4 Stunden und wartet bei jedem Stoßen auf die Dauer des Timeouts (12 Stunden).
- 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
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.