【问题标题】:Processing distributed dask collections with external code使用外部代码处理分布式 dask 集合
【发布时间】:2017-07-12 03:02:08
【问题描述】:

我已将输入数据存储为 S3 上的单个大文件。 我希望 Dask 自动切分文件,分发给工作人员并管理数据流。因此使用分布式收集的想法,例如包。

在每个工作人员上,我都有一个命令行工具 (Java),可以从文件中读取数据。因此,我想将一整块数据写入文件,调用外部 CLI/代码来处理数据,然后从输出文件中读取结果。这看起来像是处理批量数据而不是一次记录。

解决这个问题的最佳方法是什么?是否可以在worker上将分区写入磁盘并整体处理?

PS。保留在分布式集合模型中也没有必要但可取,因为对数据的其他操作可能是更简单的 Python 函数,逐条处理数据。

【问题讨论】:

标签: dask dask-distributed


【解决方案1】:

您可能需要read_bytes 函数。这会将文件分成许多块,由分隔符(如结束线)干净地分割。它会返回指向这些字节块的dask.delayed 对象列表。

此文档页面上有更多信息:http://dask.pydata.org/en/latest/bytes.html

以下是文档字符串中的示例:

>>> sample, blocks = read_bytes('s3://bucket/2015-*-*.csv', delimiter=b'\n')  

【讨论】:

  • read_bytes() 在我的嫌疑人名单上。但我有一个关于记录分隔符的问题。我显然想指定要由工作人员读取的近似块大小(例如 20MB),同时指定一个分隔符,因为记录长度会有所不同。框架将如何找出确切的分隔符位置?调度程序会读取整个文件(不受欢迎)吗?如果文件只是被切成规则的碎片,那么一些记录会被切成两半。在这种情况下,工作人员需要知道从不同的(“早期”)索引中读取?
  • read_bytes 函数根据块大小寻找位置,然后向前读取,直到找到分隔符。假设您的分隔符经常间隔,这将是相当有效的,将尊重您的大致块大小,并且将始终以分隔符结束并在一个分隔符之后开始。
  • 我用 read_bytes 做了一些实验,但是它工作正常。查看 API 我怀疑我可能会通过在数据帧上调用 map_partitions() 来获得类似/相同的效果,在从每个分区中提取的数据并返回修改后的数据帧。这样我可能能够保持在 dask dataframe API 的限制范围内。这是正确的还是我错过了什么? (为了澄清,我在这里假设我的输入数据是一个 CSV 文件)
  • 是的,如果您的文件是一个大的 csv 文件,那么 dd.read_csv 是一个很好的选择
猜你喜欢
  • 2021-09-19
  • 1970-01-01
  • 1970-01-01
  • 2018-03-26
  • 1970-01-01
  • 1970-01-01
  • 2017-01-18
  • 2020-08-19
  • 2018-01-11
相关资源
最近更新 更多