Confuso sui parametri BaseSensorOperator di Airflow: timeout, poke_interval e mode
Ho un po 'di confusione sul modo in cui funzionano BaseSensorOperatori parametri di: timeout& poke_interval. Considera questo utilizzo del sensore:
BaseSensorOperator(
soft_fail=True,
poke_interval = 4*60*60, # Poke every 4 hours
timeout = 12*60*60, # Timeout after 12 hours
)
La documentazione menziona gli atti di timeout per impostare l'attività in modo che "fallisca" dopo che è terminata. Ma sto usando un soft_fail=True, non credo che mantenga lo stesso comportamento, perché ho scoperto che l'attività non è riuscita invece di saltare dopo aver utilizzato entrambi i parametri soft_faile timeout.
Allora cosa succede qui?
- Il sensore spinge ogni 4 ore, e ad ogni colpo aspetterà la durata del timeout (12 ore)?
- O colpisce ogni 4 ore, per un totale di 3 colpi, poi va in timeout?
- Inoltre, cosa succede con questi parametri se utilizzo mode = "reschedule"?
Ecco la documentazione di 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
"""
Risposte
Definizione dei termini
poke_interval: la durata b / n successivi 'colpi' (valutazione della condizione necessaria che si sta 'percependo')timeout: Il solo colpire indefinitamente è inammissibile (se ad esempio il tuo codice buggy punta il giorno per diventare 29 ogni volta che il mese è 2, continuerà a premere fino a 4 anni). Quindi definiamo un periodo massimo oltre il quale smettiamo di premere e terminiamo (il sensore è contrassegnatoFAILEDoSKIPPED)soft_fail: Normalmente (quandosoft_fail=False), il sensore è contrassegnato comeFAILEDdopo il timeout. Quandosoft_fail=True, il sensore verrà invece contrassegnato comeSKIPPEDdopo il timeoutmode: Questo è un po 'complesso- Qualsiasi attività (incluso il sensore) quando viene eseguita, consuma un
slotin qualche pool (defaultpool o esplicitamente specificatopool); essenzialmente significa che richiede alcune risorse. - Per i sensori, questo è
- dispendioso : poiché uno slot viene consumato anche quando stiamo solo aspettando (non facendo alcun lavoro effettivo
- pericoloso : se il tuo flusso di lavoro ha troppi sensori che entrano in rilevamento nello stesso momento, possono congelare molte risorse per un bel po '. In effetti, troppe
ExternalTaskSensors sono note per mettere interi flussi di lavoro (DAG) in deadlock
- Per superare questo problema, Airflow v1.10.2 ha introdotto
modes nei sensorimode='poke'(predefinito) indica il comportamento esistente di cui abbiamo discusso sopramode='reschedule'significa che dopo un tentativo di poke , invece di andare a dormire , il sensore si comporterà come se avesse fallito (nel tentativo corrente) e il suo stato cambierà daRUNNINGaUP_FOR_RETRY. In questo modo, rilascerà il suo slot, consentendo ad altre attività di progredire mentre attende un altro tentativo di poke
- Citando lo snippet pertinente dal codice qui
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- Per maggiori informazioni leggere la sezione Parametri sensori
- Qualsiasi attività (incluso il sensore) quando viene eseguita, consuma un
E ora rispondendo direttamente alle tue domande
Q1
- Il sensore spinge ogni 4 ore, e ad ogni colpo aspetterà la durata del timeout (12 ore)?
- O colpisce ogni 4 ore, per un totale di 3 colpi, poi va in timeout?
il punto 2. è corretto
Q2
Inoltre, cosa succede con questi parametri se utilizzo mode = "reschedule"?
Come spiegato in precedenza, ognuno di questi parametri è indipendente e l'impostazione mode='reschedule'non altera in alcun modo il loro comportamento
BaseSensorOperator(
soft_fail=True,
poke_interval = 4*60*60, # Poke every 4 hours
timeout = 12*60*60, # Timeout of 12 hours
mode = "reschedule"
)
Diciamo che i criteri non vengono soddisfatti al primo colpo. Quindi verrà eseguito di nuovo dopo 4 ore di intervallo. Ma lo slot worker verrà liberato durante l'attesa poiché stiamo utilizzando l'estensione mode="reschedule".
Questo è quello che ho capito.