【问题标题】:Luigi - Overriding Task requires/inputLuigi - 覆盖任务需要/输入
【发布时间】:2018-04-23 18:22:51
【问题描述】:

我正在使用 luigi 执行一系列任务,如下所示:

class Task1(luigi.Task):
    stuff = luigi.Parameter()

    def output(self):
        return luigi.LocalTarget('test.json')

    def run(self):
        with self.output().open('w') as f:
            f.write(stuff)


class Task2(luigi.Task):
    stuff = luigi.Parameter()

    def requires(self):
        return Task1(stuff=self.stuff)

    def output(self):
        return luigi.LocalTarget('something-else.json')

    def run(self):
        with self.output().open('w') as f:
            f.write(stuff)

当我像这样启动整个工作流程时,这完全符合预期:

luigi.build([Task2(stuff='stuff')])

使用luigi.build 时,您还可以通过显式传递参数as per this example in the documentation 来运行多个任务。

但是,在我的情况下,我还希望能够完全独立地运行 Task2 的业务逻辑,而无需它参与工作流。这适用于未实现 requiresas per this example 的任务。

我的问题是,如何将这个方法作为工作流程的一部分运行,或者单独运行?显然,我可以只添加一个新的私有方法,如_my_custom_run,它获取数据并返回结果,然后在run 中使用此方法,但感觉就像是应该融入框架的东西,所以它让我觉得我误解了 Luigi 的最佳实践(仍在学习框架)。感谢您的任何建议,谢谢!

【问题讨论】:

    标签: python pipeline luigi


    【解决方案1】:

    听起来您想要dynamic requirements. 使用该示例中显示的模式,您可以读取配置或传递带有任意数据的参数,而yield 仅根据配置。

    # tasks.py
    import luigi
    import json
    import time
    
    
    class Parameterizer(luigi.Task):
        params = luigi.Parameter() # Arbitrary JSON
    
        def output(self):
            return luigi.LocalTarget('./config.json')
    
        def run(self):
            with self.output().open('w') as f:
                json.dump(params, f)
    
    class Task1(luigi.Task):
        stuff = luigi.Parameter()
    
        def output(self):
            return luigi.LocalTarget('{}'.format(self.stuff[:6]))
    
        def run(self):
            with self.output().open('w') as f:
                f.write(self.stuff)
    
    
    class Task2(luigi.Task):
        stuff = luigi.Parameter()
        params = luigi.Parameter()
    
    
        def output(self):
            return luigi.LocalTarget('{}'.format(self.stuff[6:]))
    
        def run(self):
    
            config = Parameterizer(params=self.params)
            yield config
    
            with config.output().open() as f:
                parameters = json.load(f)
    
            if parameters["runTask1"]:
                yield Task1(stuff=self.stuff)
            else:
                pass
            with self.output().open('w') as f:
                f.write(self.stuff)
    
    if __name__ == '__main__':
        cf_json = '{"runTask1": True}'
    
        print("Trying to run with Task1...")
        luigi.build([Task2(stuff="Task 1Task 2", params='{"runTask1":true}')], local_scheduler=True)
    
        time.sleep(10)
    
        cf_json = '{"runTask1": False}'
    
        print("Trying to run WITHOUT Task1...")
        luigi.build([Task2(stuff="Task 1Did just task 2", params='{"runTask1":false}')], local_scheduler=True)
    

    (只需调用python tasks.py即可执行)

    我们可以很容易地想象将多个参数映射到多个任务,或者在允许执行各种任务之前应用自定义测试。我们也可以重写它以从luigi.Config 获取参数。

    还要注意来自Task2 的以下控制流:

        if parameters["runTask1"]:
            yield Task1(stuff=self.stuff)
        else:
            pass
    

    在这里,我们可以运行替代任务,或者动态调用任务,正如我们在 luigi 存储库的示例中看到的那样。例如:

        if parameters["runTask1"]:
            yield Task1(stuff=self.stuff)
        else:
            # self.stuff is not automatically parsed to int, so this list comp is valid
            data_dependent_deps = [Task1(stuff=x) for x in self.stuff] 
            yield data_dependent_deps
    

    这可能比简单的run_standalone() 方法涉及更多,但我认为它最接近您在记录的 luigi 模式中寻找的内容。

    来源:https://luigi.readthedocs.io/en/stable/tasks.html?highlight=dynamic#dynamic-dependencies

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2022-08-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-06-12
      相关资源
      最近更新 更多