如何使用气流检查长时间运行的http任务的状态?

sau*_*abh 1 job-scheduling airflow

我的用例是使用气流控制跨微服务的大量预定作业。我正在尝试的解决方案是使用气流作为集中式作业调度程序并通过进行 http 调用来触发作业。其中一些作业将运行很长时间,例如。超过 10 分钟或长达 1 小时。

如何通过气流定期检查这些作业的状态?如果远程任务已完成但气流不知道作业成功怎么办?我可以将作业完成事件发布到 kafka 并让气流在 kafka 上监听以获取作业状态吗?

Mik*_*ike 6

您可以通过多种方式使用 Airflow 和您的微服务来做到这一点。通常,您会想要使用传感器,这是适合此类情况的 Airflow 对象。首先查看BaseSensorOperator和有关操作符。在 Airflow 中,Sensor 就像 Operator 一样使用(传感器就是 Operator)。所以你可以创建一个这样的工作:

http_post_task -> http_sensor_task -> success_task
Run Code Online (Sandbox Code Playgroud)

http_post_task 将触发作业,http_sensor_task 将定期检查以查看作业是否完成(例如 GET 请求微服务并检查 200,也许?),并且 success_task 将在 http_sensor_task 成功后执行。

您的 http_sensor_task 需要是您自己的自定义传感器。这里有一些 sudo 代码可以帮助你创建这个传感器(记住传感器就像操作员一样使用)。考虑您向微服务发出请求,然后发出另一个请求以检查作业状态(GET 请求并检查 200)的情况,您将像这样扩展 BaseSensorOperator:

from airflow.operators.sensors import BaseSensorOperator
from airflow.utils.decorators import apply_defaults
from time import sleep
import requests

class HTTPSensorOperator(BaseSensorOperator): 
    """
    Pokes a URL until it returns 200
    """
    ui_color = '#000000'
    @apply_defaults
    def __init__( self, url, *args, **kwargs):
        super(HTTPSensorOperator, self).__init__(*args, **kwargs)
        self.url = url


    def poke(self, context):
        """
        GET request url and return True if response is 200, False otherwise
        """
        r = requests.post(self.url)
        if r.status_code == 200:
            return True
        else:
            return False

    def execute(self, context):
        """
        Check the url and wait for it to return 200.
        """
        started_at = datetime.utcnow()
        while not self.poke(context):
            if (datetime.utcnow() - started_at).total_seconds() > self.timeout:
                if self.soft_fail:
                    raise AirflowSkipException("Exporting {0}/{1} took to long.".format(self.project, self.instance))
                else:
                    raise AirflowSkipException("Exporting {0}/{1} took to long.".format(self.project, self.instance))
            sleep(self.poke_interval)
        self.log.info("Success criteria met. Exiting.")
Run Code Online (Sandbox Code Playgroud)

然后使用运算符,如:

http_sensor_task = HTTPSensorOperator(
      task_id="http_sensor_task",
      url="http://localhost/check_job?job_id=1",
      timeout=3600, # 1 hour
      dag=dag
   )
Run Code Online (Sandbox Code Playgroud)

因此,您必须决定您的微服务将如何与 Airflow 通信。就在我的脑海中,我想你会提出 1 个请求来触发一项工作,然后提出后续请求(可能需要 10 秒)来检查一项工作。祝你好运!