minhajuddin

minhajuddin

Like the title says, I want to pass a stream to two outputs.

I am looking for something like below:

function_which_returns_stream
|> Enum.tee( mapper_for_func1 |> postgrex_stream1, mapper_for_func2 |> postgrex_stream2)

Is there a way to achieve this in Elixir?

Showing Posts 1 to 10

wmnnd

wmnnd

You could simply use Enum.map/2 with an anonymous function that calls your other two functions like this:


fn1 = fn (x) -> IO.puts("1: #{x}") end
fn2 = fn (x) -> IO.puts("2: #{x}") end
1..3
|> Stream.map(&(&1))
|> Enum.map(stream, fn (x) ->
     fn1.(x)
     fn2.(x)
   end)
pma

pma

If you’re streaming to external systems (postgres), you can also use GenStage and a Broadcast Dispatcher? This would give you concurrency and back-pressure.

minhajuddin

minhajuddin OP

Genstage seems too heavy for my use case :slight_smile: If there is no alternative, I’ll probably have to use it.

minhajuddin

minhajuddin OP

Thanks. This really doesn’t stream the data but calls the other 2 functions for each element. I need to stream almost 2million rows to a postgres stream.

benwilson512

benwilson512

Author of Craft GraphQL APIs in Elixir with Absinthe

Can you elaborate a bit on your actual scenario?

Qqwy

Qqwy

TypeCheck Core Team

What about something like this?

defmodule Test do
  def tee(input_stream, functions) do
    Stream.map(functions, fn fun ->
      run_async_pipeline(input_stream, fun)
    end)
    |> Stream.concat
    |> Stream.run
  end

  def run_async_pipeline(stream, fun) do
    stream
    |> Task.async_stream(fun)
    |> Stream.map(fn {:ok, item} -> item end) # error handling could be added here.
    |> put_in_appropriate_postgres_stream
  end
end

tee([1,2,3,4], [&(&1*&1), &(&1+1), &(&1-1)])
minhajuddin

minhajuddin OP

Thanks, this is very close to what I want. However, this seems to reread the stream. I tried it with the following code:

defmodule Test do
  def tee(input_stream, functions) do
    Stream.map(functions, fn fun ->
      run_async_pipeline(input_stream, fun)
    end)
    |> Stream.concat
    |> Stream.run
  end

  def run_async_pipeline(stream, fun) do
    stream
    |> Task.async_stream(fun)
    |> Stream.map(fn {:ok, item} -> item end) # error handling could be added here.
    |> Enum.to_list
    |> IO.inspect(label: "FUN")
  end
end

stream = Stream.unfold(5, fn 0 -> nil; n -> IO.puts("#{n}."); {n, n-1} end)
Test.tee(stream, [&(&1*&1), &(&1+1), &(&1-1)])

And got this output

5.
4.
3.
2.
1.
FUN: [25, 16, 9, 4, 1]
5.
4.
3.
2.
1.
FUN: [6, 5, 4, 3, 2]
5.
4.
3.
2.
1.
FUN: [4, 3, 2, 1, 0]

My use case parses some xml files and pipes them to two postgresql table streams. I don’t want to parse the xml file 2 times to run this. Is there a way to do that?

minhajuddin

minhajuddin OP

This is the high level code that I have:

defmodule Loader do
  import Logger, only: [debug: 1]
  import PG

  def load_all do
    files()
    |> Enum.with_index(1)
    |> Task.async_stream(&load_file/1, max_concurrency: DataLoader.jobs, timeout: :infinity)
    |> Stream.flat_map(fn {:ok, products} -> products end)
    |> Stream.uniq_by(fn [id | _] -> id  end)
    |> copy("COPY products (id, name) FROM STDIN DELIMITERS E'\\t' NULL ''")
  end

  defp load_file({filepath, index}) do
    debug "#{index} loading #{filepath}"
    Parser.parse_xml(filepath)
  end

  defp files do
    Path.wildcard("#{DataLoader.data_dir}/products_*.xml")
  end
end

The current code streams the data to a postgresql stream to a single table. However, I get info about the product as well as the product_categories in the same input file. I need to populate both the tables using these input files. I don’t want to do 2 passes for parsing.

bbense

bbense

My 2 cents:

I would start a separate process that streamed incoming messages into the second stream.
Then have the initial stream output those messages to the second stream.

Genstage is a much fancier version of this that handles all the horrible edge cases.

minhajuddin

minhajuddin OP

Thanks, seems like Genstage is the easy way out for this. Would be a lot easier if there was a simpler version :slight_smile:

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
nseaSeb
Hello, I know there is an approach for handling lists that allows for optimized traversal, but I can’t recall the specific method (somet...
New
brecabral
Documentation While reading the Scoped Routes section, I noticed that the documentation currently refers to a problem without explainin...
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
asweet-confluent
I recently noticed that Elixir’s Logger defaults its primary log level to :debug when no :logger, :level application configuration is pre...
New

Other Trending Topics Top

JesseHerrick
Hey, I’m Jesse and I’m the main contributor behind Dexter, a full-featured, lightning-fast Elixir LSP optimized for large codebases. It s...
New
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
mhanberg
Hi everyone! The first release candidate for the Expert language server project is now available! We’ve published a press release detai...
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

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews