【问题标题】:switching from luigi to airflow从 luigi 切换到气流
【发布时间】:2017-07-11 05:47:18
【问题描述】:

我有一个相对简单的任务,首先在 1.2 个 mio 文件上运行并为每个文件创建一个管道(其中包含保存中间产品的多个步骤)。我已经在 luigi 中实现了这一点:https://gist.github.com/wkerzendorf/395c85a2955002412be302d708329f7f。我喜欢 Luigi 使用文件系统来查看任务是否已完成。 我还找到了一个实现,我可以删除一个中间产品,管道将重新创建所有依赖产品(所以我可以更改管道)。 我将如何在气流中做到这一点(或者我应该坚持使用 Luigi?)?

【问题讨论】:

    标签: airflow luigi


    【解决方案1】:

    我真的不知道 Luigi 是如何工作的。我主要使用 Apache Airflow。 Airflow 是一个工作流管理系统。这意味着它不会传输数据、转换数据或生成一些数据(尽管它会生成日志并且有一个名为Xcom 的概念允许在任务之间交换消息,从而允许更细微的控制形式和共享状态。),例如。阿帕奇尼菲。但它定义了您使用Operators 实例化它的每个任务的依赖关系,例如。 BashOperator。为了知道某项任务是否完成,它会检查同一任务返回的信号。

    以下是您希望在 Airflow 中实现的示例。

    from airflow.operators.bash_operator import BashOperator
    from airflow.operators.python_operator import PythonOperator
    import glob
    import gzip
    import shutil
    
    args = {
        'owner': 'airflow',
        'start_date': airflow.utils.dates.days_ago(2)
    }
    
    dag = DAG(
        dag_id='example_dag', default_args=args,
        schedule_interval='0 0 * * *',
        dagrun_timeout=timedelta(minutes=60))
    
    
    def extract_gzs():
        for filename in glob.glob('/1002/*.gz')
            with gzip.open(filename, 'rb') as f_in, open(filename[:-3], 'wb') as f_out:
                shutil.copyfileobj(f_in, f_out)
    
    
    extractGZ = PythonOperator(
        task_id='extract_gz',
        provide_context=True,
        python_callable=extract_gzs(),
    dag=dag)
    
    
    cmd_cmd="""
    your sed script!
    """
    
    sed_script = BashOperator(
        task_id='sed_script', 
        bash_command=cmd_cmd, 
        dag=dag)
    
    
    extractGZ.set_downstream(sed_script)
    
    1. 导入您想在 Airflow 中使用的运算符(当然,如果您需要其他类/库)
    2. 定义你的 Dag。在变量args 中,我定义了ownerstart_date 参数。
    3. 然后实例化您的 DAG。在这里,我将其命名为 example_dag,将其定义变量归为 schedule_interval 以及超时时间(根据您的需要还有更多参数可供使用)
    4. 创建了一个python函数extract_gzs()
    5. 实例化了一个PythonOperator,我在其中调用了我的python func
    6. 对 bash 代码做同样的事情
    7. 确定两个任务实例之间的依赖关系

    当然,还有更多方法可以实现相同的想法。根据需要进行调整! PS:Here 有一些 Apache Airflow 的例子

    【讨论】:

    • 也许,我误解了整个管道的事情。我假设您构建了一个管道,然后向它提供一个数据集,然后沿着该管道的线进行转换。这将意味着一个适用于 1.2 mio 文件的管道。这不是思考这个问题的正确方法吗?制作一个在文件上运行 sed 的气流管道,然后将其应用于 1.2 mio 文件应该是微不足道的,不是吗?
    • @WolfgangKerzendorf 检查我修改后的答案。
    • 谢谢,我会试一试。它仍然不完全是我想象的那样。我想为一个文件构建一个管道,然后以某种方式通过这个东西推动每个文件。
    猜你喜欢
    • 2023-03-25
    • 2017-07-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多