fchabouis

fchabouis

I have a zip file, containing CSVs, on a remote S3 cellar.

Using Unzip, I can get a stream of the CSV file that I would like to decode, like so :

    aws_s3_config =
      ExAws.Config.new(:s3,
        access_key_id: ["xxx", :instance_role],
        secret_access_key: ["xxx", :instance_role]
      )

    file = new(zip_name, bucket_name, aws_s3_config)
    {:ok, unzip} = Unzip.new(file)
    stream = Unzip.file_stream!(unzip, file_name)

as explained in the doc.

Now I would like to consume that stream by reading it with CSV.
So I try stream |> CSV.decode |> Enum.take(1) and get an error ** (FunctionClauseError) no function clause matching in CSV.Decoding.Preprocessing.Lines.starts_sequence?/5

If I write the content of my CSV on the disk and then read it, it works fine :

# write the file on disk
stream |> Stream.into(File.stream!("stops.txt")) |> Stream.run()
# then read and decode it
File.stream!("stops.txt") |> CSV.decode |> Enum.take(1)

I get the desired result, the first row of the CSV file : [ok: ["\uFEFFstop_id", "stop_name", "stop_lat", "stop_lon", "location_type"]]

The difference I see is that Unzip.file_stream! and File.stream!("stops.txt") do not stream the file the same way. Unzip seem to do it by chunks of 65k, while File.stream! streams line by line.

How can I solve this, without writing the file to disk as an intermediary step?
Thanks!

Showing Posts 1 to 8

ahamez

ahamez

Hello,

I don’t know if it will help, as I’m not using Unzip, but StreamGzip, in combination with NimbleCSV, for this purpose. But maybe it will give you some hint?

I have the following function that returns a stream for an object downloaded from S3:

  defp get_object_stream(object) do
    {:ok, io_pid} = StringIO.open(object)

    io_pid
    |> IO.binstream(4096)
    |> StreamGzip.gunzip()
    |> NimbleCSV.RFC4180.to_line_stream()
  end

In my case, the trick was to use to_line_stream.

I then can use this stream like this:

object
|> get_object_stream()
|> NimbleCSV.RFC4180.parse_stream()

As you can see, I’m not streaming directly from S3 as I first download the object in memory, but if you have something that’s already able to stream from S3, you would just have to replace the part that constructs the stream from the in-memory string with your stream from S3.

akash-akya

akash-akya

Echoing what @ahamez has already mentioned, the issue seems to be that CSV.decode expects stream of lines. But Unzip.file_stream! returns stream of blobs. You can convert stream of blobs to stream of lines yourself, or you can use NimbleCSV as already mentioned.

Unzip.file_stream!(unzip, file_name)
|> NimbleCSV.RFC4180.to_line_stream()
|> NimbleCSV.RFC4180.parse_stream()
fchabouis

fchabouis OP

Thanks for you replies.

Unfortunately, applying to_line_stream and parse_stream yields an error:

stream 
|> NimbleCSV.RFC4180.to_line_stream() 
|> NimbleCSV.RFC4180.parse_stream() 
|> Stream.run()
** (ArgumentError) errors were found at the given arguments:

  * 1st argument: not a bitstring
    :erlang.bit_size(["\uFEFFservice_id,monday,tuesday,wednesday,thursday,friday,saturday,sunday,start_date,end_date\r\nE1-5-1-127,1,1,1,1,1,1,1,20220115,20220115\r\nH2-0-1-1,1,0,0,0,0,0,0,20220103,20220415\r\nH2-0 (...)
    (nimble_csv 1.2.0) lib/nimble_csv.ex:393: NimbleCSV.RFC4180.to_line_stream_chunk_fun/3
    (elixir 1.12.2) lib/stream.ex:264: anonymous fn/4 in Stream.chunk_while_fun/2
    (elixir 1.12.2) lib/enum.ex:4280: Enumerable.List.reduce/3
    (elixir 1.12.2) lib/stream.ex:931: Stream.do_list_transform/7
    (elixir 1.12.2) lib/stream.ex:1719: Enumerable.Stream.do_each/4
    (elixir 1.12.2) lib/stream.ex:880: Stream.do_transform/5
    (elixir 1.12.2) lib/stream.ex:649: Stream.run/1

What I don’t understand, is the return structure of Unzip.file_stream!. I would expect the stream to return some data, chunk by chunk. If I do Unzip.file_stream! |> Enum.to_list(), I thought I would get something like ["some binary data", "some other binary data", "..."]. Instead I get a nested list of data, that looks like this :

[
 [
  [
   ["some data"], 
   "some other data"], 
   "data again"`]
]
LostKobrakai

LostKobrakai

That’s why to_line_stream fails. It expects an enumerable of binaries.

fchabouis

fchabouis OP

Ok thanks, I needed to adapt the stream coming from Unzip like so :

Unzip.file_stream!(unzip, file_name) 
|> Stream.map(fn c -> List.flatten(c) |> Enum.join("") end) 
|> NimbleCSV.RFC4180.to_line_stream() 
|> NimbleCSV.RFC4180.parse_stream() 
|> Enum.to_list() 

Thank you all for your kind assistance and for giving me good tips!

LostKobrakai

LostKobrakai

List.flatten(c) |> Enum.join("") would probably better replaced with IO.iodata_to_binary/1

fchabouis

fchabouis OP

Nice, thanks @LostKobrakai :+1:

fchabouis

fchabouis OP

I wrote a blog post on the subject, if than can be useful to someone.

— All posts loaded —

Where Next? Top

Trending in Questions Top

RSP87
I’m working on a project that simulates the bumbl example in the programming phoenix book. It acts almost like an email client. We have a...
New
kszambelanczyk
Hello! Could someone please give me a help/sample code, how to delete a file from s3 using waffle/waffle_ecto from Phoenix app. I creat...
New
RemyXRenard
I’m seeing that a list inside a Kino.DataTable will be interpreted as a charlist, even if the Kino.configure() is set to charlists: :as_l...
New
velrest
So my question is quite simple and i have found no conclusive answer on forum, google or AI. Should we use :erlang.float for Integer to ...
New
samoloth
Hi, I’ve just set up an application with ash_authentication. There is only magic link strategy for now, so there is no confirmation add o...
New
FlyingNoodle
If a change or preparation module uses Ash.Changeset.get_argument/2 or Ash.Query.get_argument/2 (or any of the other get_argument functio...
New
ryanwinchester
apply_graft/2 doesn’t rewrite an add_many sub-workflow’s deps on an add step. Grafted jobs cancel with “upstream job was deleted” Version...
New

Other Trending Topics Top

mudasobwa
I am happy to introduce the very α version of the new programming language compiled to BEAM. Welcome Cure. It has literally three kille...
New
marciok
Hi there! We created Gust: A task orchestrator inspired by Airflow. For those who have never heard about Aiflow, it’s a Python-based wor...
New
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Dmk
Xamal is a deployment tool for Elixir apps that deploys native releases to bare metal servers over SSH. It’s a port of GitHub - basecamp/...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge & Solve. They are GUI (Emerge) and State management (S...
New
netoum
Corex is an accessible, unstyled UI component library for Phoenix that integrates Zag.js state machines using Vanilla JavaScript and Live...
New

Latest on Elixir Forum

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews