【问题标题】:luigi upstream task should run once to create input for set of downstream tasksluigi 上游任务应该运行一次以为下游任务集创建输入
【发布时间】:2017-05-10 21:36:23
【问题描述】:

我有一个很好的直接工作管道,我在命令行上通过 luigi 运行的任务会触发所有必需的上游数据获取并以正确的顺序进行处理,直到它流入我的数据库。

class IMAP_Fetch(luigi.Task):
  """fetch a bunch of email messages with data in them"""
  date = luigi.DateParameter()
  uid = luigi.Parameter()
…
  def output(self):
    loc = os.path.join(self.data_drop, str(self.date))
    # target for requested message
    yield LocalTarget(os.path.join(loc, uid+".msg"))

  def run(self):
     # code to connect to IMAP server and run FETCH on given UID
     # message gets written to self.output()
…

class RecordData(luigi.contrib.postgres.CopyToTable):
  """copy the data in one email message to the database table"""
  uid = luigi.Parameter()
  date = luigi.DateParameter()
  table = 'msg_data'
  columns = [(id, int), …]

  def requires(self):
    # a task (not shown) that extracts data from one message  
    # which in turn requires the IMAP_Fetch to pull down the message
    return MsgData(self.date, self.uid) 

  def rows(self):
    # code to read self.input() and yield lists of data values 

好东西。不幸的是,第一次数据获取与远程 IMAP 服务器通信,每次获取都是一个新连接和一个新查询:非常慢。我知道如何在一个会话(任务实例)中获取所有单独的消息文件。我不明白如何让下游任务保持原样,一次处理一条消息,因为需要一条消息的任务会触发仅获取一条消息,而不是获取所有可用消息。我为错过明显的解决方案而提前道歉,但到目前为止,它让我很难过如何保持我漂亮的简单愚蠢的管道基本上保持原样,但让顶部的漏斗在一次调用中吸收所有数据。感谢您的帮助。

【问题讨论】:

  • 请附上一些代码,说明您的任务是做什么的以及我们如何为您提供帮助
  • 按要求更新了我原来的帖子

标签: luigi data-pipeline


【解决方案1】:

您的解释中我缺少的是发送到RecordData 任务的uid 值列表来自哪里。对于这个解释,我假设您有一组 uid 值,您希望将它们合并到一个 ImapFetch 请求中。

一种可能的方法是定义一个batch_id,另外定义一个uid,其中batch_id 是指您希望在单个会话中获取的消息组。其中@987654328 之间的关联@ 和 batch_id 的存储取决于您。它可以是传递给管道的参数,也可以是外部存储的。您遗漏的任务MsgData,其requires 方法目前返回一个带有uid 参数的ImapFetch 任务,应该改为需要一个带有batch_id 参数的ImapFetch 任务。 MsgData 任务所需的第一个ImapFetch 任务将检索与该batch_id 关联的所有uid 值,然后在单个会话中检索这些消息。所有其他MsgData 任务都需要(并正在等待)这一批ImapFetch 才能完成,然后它们都可以像管道的其余部分一样执行各自的消息。因此,调整批量大小可能对整体处理吞吐量很重要。

另一个缺点是它在批次级别而不是单个项目级别的原子性较低,因为如果没有成功检索到 uid 值中的一个,则批次 ImapFetch 将失败。

第二种方法是将 Imap 会话作为每个进程(工作人员)更多的单一资源打开,并让 ImapFetch 任务重用相同的会话。

【讨论】:

  • 非常感谢您清晰周到的解释。好吧,在我发布这个之后,你的第一个提案也给了我,所以你已经为我验证了它。我也喜欢第二个想法。再次感谢。
【解决方案2】:

你也可以使用这样的计数器。

class WrapperTask(luigi.WrapperTask):
      counter = 0 
      def requires(self):
         if self.counter == 0:  # check if this is the first time the long DB process is called
         ## do your process here. This executes only once and is skipped next time due to the counter
         self.counter = self.counter + 1
         return OtherTask(parameters)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-09-07
    • 1970-01-01
    • 1970-01-01
    • 2020-01-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-08
    相关资源
    最近更新 更多