【问题标题】:Can I use the output of one DynamicResource after a map() for multiple solids?我可以在 map() 之后为多个实体使用一个 DynamicResource 的输出吗?
【发布时间】:2021-09-28 10:06:52
【问题描述】:

我正在做类似于文档中的dynamic mapping and collect 示例。该示例列出了目录中的文件,将每个文件映射到计算文件大小的实体,然后收集输出以汇总总体大小。

但是,我想在每个实体上并行运行多个实体。所以继续这个例子:我会列出目录中的文件;然后映射,以便为每个文件计算大小,检查文件权限,并并行计算 md5sum;最后收集输出。

我可以在每个文件上按顺序运行这些,例如:

file_results = list_files()
.map(compute_size)
.map(check_permissions)
.map(compute_md5sum)
summarize(file_results.collect())

但如果这些实际上不是串行依赖项,那么并行处理每个文件的工作会很好。有没有这样的语法:

file_results = list_files().map(
compute_md5sum(check_permissions(compute_size)))
summarize(file_results.collect())

【问题讨论】:

    标签: dagster


    【解决方案1】:

    如果我理解正确,这样的事情应该可以完成您正在寻找的内容:

    def _process_file(file):
        size = compute_size(file)
        perms = check_permissions(file)
        hash = compute_md5sum(file)
        return summarize_file(size, perms, hash)
    
    
    file_results = list_files().map(_process_file)
    summarize(file_results.collect)
    

    【讨论】:

    • 谢谢,成功了!您介意澄清一下您的代码 sn-p 应该在 @pipeline 方法中吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-12-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-07-19
    相关资源
    最近更新 更多