【问题标题】:How do you write a file per chunk in a Stream in Elixir如何在 Elixir 的 Stream 中为每个块写入文件
【发布时间】:2017-06-15 04:18:40
【问题描述】:

我有一个问题,我需要读取一个非常大的文件,然后打印每个块的解析结果。最后没有一个完整的列表。

到目前为止,我可以在 MapSet 中获得 uniq 结果,但无法弄清楚如何根据 chunk_size 写入文件

使用此方法获取唯一的文件名

def new_file_name do
  hex = :crypto.hash(:md5, Integer.to_string(:os.system_time(:millisecond)))
    |> Base.encode16
end 

到目前为止,我拥有的最好的是这个,它给了我一个 MapSet 列表,其中包含块大小的独特结果。这是一个 MapSet 列表,最终可能会因内存太大而无法容纳。

def parse(file_path, chunk_size) do
  file_path
    |> File.stream!(read_ahead: chunk_size)
    |> Stream.drop(1)  # remove header
    |> Stream.map(&"#{&1}\")  # Prepare to be written as a csv
    |> Stream.chunk(chunk_size, chunk_size, [])  # break up into chunks
    |> method # method to write per chunk to file. 
end

我之前有的是

|> Stream.map(&MapSet.new(&1))  # Create MapSet of unique values from each chunk

但我似乎找不到任何将 MapSet 写入文件的示例。

【问题讨论】:

  • 在调用Enum 函数或Stream.run/1 之一之前不会进行任何计算。所以你可能想完成使用Enum.map 而不是Stream.map
  • 从文档hexdocs.pm/elixir/MapSet.html看来,您唯一能做的就是将其转换为列表(使用Mapset.to_list(map_set) ),然后将其写入文件。我自己没有尝试过 - 所以如果您想稍后再读回数据,可能需要特别小心。
  • MapSet 由什么组成?您想保留在新文件中的行吗?如果您放弃在MapSet 中存储哪些数据的示例,这对您来说会是个问题吗?这将有助于理解问题和建议方法
  • @GavinBrelstaff 这是我尝试过的。帕维尔->是的。它只是一列字符串。没什么特别的。只是一个唯一字符串的 MapSet。我喜欢 MapSet.to_list 但如何为列表中的每个列表编写一个文件?
  • 你能发布一个示例输入和输出吗?您想如何将 MapSet 中的元素写入文件?

标签: elixir


【解决方案1】:

您可以使用Enum.reduce/3 和文件句柄作为累加器来打开文件一次,然后一次写入一个块:

def parse(file_path, chunk_size) do
  file_path
  |> File.stream!(read_ahead: chunk_size)
  |> Stream.drop(1)  # remove header
  |> Stream.map(&"#{&1}\")  # Prepare to be written as a csv
  |> Stream.chunk(chunk_size, chunk_size, [])  # break up into chunks
  |> Enum.reduce(File.open!("output.txt", [:write]), fn chunk, file ->
    :ok = IO.write(file, chunk)
    file
  end)
end

您可能想要调整将块写入文件的方式。以上将chunk视为iodata,有效地连接块中的字符串并写入。

如果您只想为每个块编写唯一项目,您可以添加:

|> Stream.map(fn chunk -> chunk |> MapSet.new |> MapSet.to_list end)

在输入Enum.reduce/3之前。

【讨论】:

  • 使用|> Stream.map 版本效果很好。只需要以Stream.run结束即可
  • 你在哪里添加Stream.run
  • 在流的末尾所以|> Stream.map(**Stuff**) |> Stream.map(Enum reduce method per chunk to file) |> Stream.run
  • 啊,我以为你想将所有块写入同一个文件,在这种情况下,Enum.reduce 将是顶层的最后一个管道,你不需要Stream.run
【解决方案2】:

在@Dogbert 的帮助下找到了一种有趣的方法。使用 Stream 会将我锁定为最大 100% 的 cpu 使用率。有了这个,我能够达到最高 256% 的 cpu 使用率。这是在每个 300MB 的几个文件上运行的。 30分钟解析。

def alt_flow_parse_dir(path, out_file, chunk_size) do
  concat_unique =  File.open!(path <> "/" <> out_file, [:read, :utf8, :write])

  Path.wildcard(path <> "/*.csv")
    |> Flow.from_enumerable
    |> Flow.map(&append_to_file(&1, path, concat_unique, chunk_size))
    |> Flow.run

  File.close(concat_unique)
end

# I just want the unique items of the first column
def append_to_file(filename, path, out_file, chunk_size) do
  file = filename
    |> String.split("/")
    |> Enum.take(-1)
    |> List.to_string
  path <> file
    |> File.stream!
    |> Stream.drop(1)
    |> Flow.from_enumerable
    |> Flow.map(&String.split(&1, ",") |> List.first)
    |> Flow.map(&String.trim(&1,"\n"))
    |> Flow.partition
    |> Stream.chunk(chunk_size, chunk_size, [])
    |> Flow.from_enumerable
    |> Flow.map(fn chunk ->
        chunk
          |> MapSet.new
          |> MapSet.to_list
          |> List.flatten
      end)
    |> Flow.map(fn line ->
        Enum.map(line, fn item ->
            IO.puts(out_file, item)
          end)
        end)
     |> Flow.run
  end

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-02-04
    • 2012-11-16
    • 1970-01-01
    • 2021-11-07
    • 2012-09-02
    • 2011-03-23
    • 1970-01-01
    • 2015-12-11
    相关资源
    最近更新 更多