Confuso sui parametri BaseSensorOperator di Airflow: timeout, poke_interval e mode

Sep 07 2020

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?

  1. Il sensore spinge ogni 4 ore, e ad ogni colpo aspetterà la durata del timeout (12 ore)?
  2. O colpisce ogni 4 ore, per un totale di 3 colpi, poi va in timeout?
  3. 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

6 y2k-shubham Sep 07 2020 at 11:44

Definizione dei termini

  1. poke_interval: la durata b / n successivi 'colpi' (valutazione della condizione necessaria che si sta 'percependo')

  2. 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 è contrassegnato FAILEDo SKIPPED)

  3. soft_fail: Normalmente (quando soft_fail=False), il sensore è contrassegnato come FAILEDdopo il timeout. Quando soft_fail=True, il sensore verrà invece contrassegnato come SKIPPEDdopo il timeout

  4. mode: Questo è un po 'complesso

    • Qualsiasi attività (incluso il sensore) quando viene eseguita, consuma un slotin qualche pool ( defaultpool o esplicitamente specificato pool); 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 mode s nei sensori
      • mode='poke' (predefinito) indica il comportamento esistente di cui abbiamo discusso sopra
      • mode='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à da RUNNINGa UP_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

E ora rispondendo direttamente alle tue domande

Q1

  1. Il sensore spinge ogni 4 ore, e ad ogni colpo aspetterà la durata del timeout (12 ore)?
  2. 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

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"
)

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.