Maxximiliann
defmodule AsyncProcessor do
@max_concurrency System.schedulers_online() * 30
def format(sequence, :flat_map_stream),
do: Stream.flat_map(sequence, fn {:ok, result} -> result end)
def format(sequence, :stream), do: Stream.map(sequence, fn {:ok, result} -> result end)
def format(sequence, :flat_map),
do: Enum.reduce(sequence, [], fn {:ok, [_ | _] = result}, acc -> result ++ acc end)
def format(sequence, :enum_map), do: Enum.map(sequence, fn {:ok, result} -> result end)
def execute(enumerable, lambda, output_config \\ :enum_map, timeout \\ 20_000) do
Task.async_stream(enumerable, &lambda.(&1),
max_concurrency: @max_concurrency,
timeout: timeout
)
|> format(output_config)
end
end
defmodule Print do
@reset IO.ANSI.reset()
@blue IO.ANSI.color(33)
@fuscia IO.ANSI.color(161)
@light_green IO.ANSI.color(44)
@mid_green IO.ANSI.color(76)
@yellow IO.ANSI.color(226)
def utc_time, do: DateTime.utc_now()
def text(message) do
timestamp = @blue <> "#{utc_time()} UTC-" <> @mid_green
IO.puts("#{timestamp} #{message} \n")
end
def text(details, module, line_number) do
timestamp = @blue <> "#{utc_time()} UTC-" <> @mid_green
IO.puts("#{timestamp} #{module} - line #{line_number}: #{details} \n")
end
def error(details, module, line_number) do
timestamp = @blue <> "#{utc_time()} UTC-"
formatted_details = @fuscia <> "#{details}" <> @reset
IO.puts("#{timestamp} #{module} - line #{line_number}: #{formatted_details} \n")
end
def warning(details, module, line_number) do
timestamp = @blue <> "#{utc_time()} UTC-"
warning_message = @yellow <> "#{details}" <> @reset
IO.puts("#{timestamp} #{module} - line #{line_number}: #{warning_message} \n")
end
def highlight(details, module, line_number) do
timestamp = @blue <> "#{utc_time()} UTC-" <> @light_green
IO.puts("#{timestamp} #{module} - line #{line_number}: #{details} \n")
end
def highlight(message) do
timestamp = @blue <> "#{utc_time()} UTC-"
text = @light_green <> "#{message}"
IO.puts("#{timestamp} #{text} \n")
end
end
defmodule LogBook do
require Logger
def write_to_log(data, module, line_number) do
data_eval = inspect(data, limit: :infinity)
Logger.info("#{module} - line #{line_number}: #{data_eval}")
{:ok, data_eval}
end
def write_to_console(:quiet),
do: {:ok, :message_logged}
def write_to_console(_, _, _, :quiet),
do: write_to_console(:quiet)
def write_to_console(data_eval, module, line_number, :error),
do: Print.error(data_eval, module, line_number)
def write_to_console(data_eval, module, line_number, :warning),
do: Print.warning(data_eval, module, line_number)
def write_to_console(data_eval, module, line_number, message_type)
when message_type == :trace or message_type == nil,
do: Print.highlight(data_eval, module, line_number)
def write_to_console(_, :quiet),
do: write_to_console(:quiet)
def write_to_console(message, :print_to_screen),
do: Print.text(message)
def write_to_console(message, :highlight),
do: Print.highlight(message)
def main(data, module, line_number, message_type \\ :trace, list? \\ false)
def main(message, module, line_number, message_type, false)
when message_type == :print_to_screen or message_type == :highlight do
with {:ok, _} <- write_to_log(message, module, line_number) do
write_to_console(message, message_type)
else
glitch ->
raise "#{__MODULE__}: Mishandled value: #{inspect(glitch)}"
end
end
def main(data, module, line_number, message_type, false) do
with {:ok, data_eval} <- write_to_log(data, module, line_number) do
write_to_console(data_eval, module, line_number, message_type)
else
glitch ->
raise "#{__MODULE__}: Mishandled value: #{inspect(glitch)}"
end
end
def main([_ | _] = data, module, line_number, message_type, true) do
with {:ok, _data_eval} <- write_to_log(data, module, line_number) do
AsyncProcessor.execute(
data,
&(inspect(&1, limit: :infinity)
|> write_to_console(module, line_number, message_type))
)
else
glitch ->
raise "#{__MODULE__}: Mishandled value: #{inspect(glitch)}"
end
end
def testing(data, module, line_number),
do: main(data, module, line_number)
def testing(data, module, line_number, opts) do
list? = opts[:list] || false
message_type = opts[:message_type] || :trace
main(data, module, line_number, message_type, list?)
end
end
defmodule Ventris do
def update_keys_list(keys_list, acc) when is_list(acc) do
case is_list(keys_list) do
true -> keys_list ++ acc
false -> [keys_list] ++ acc
end
end
def extract_nested_keys(nested_api_data, root_key) when is_map(nested_api_data) do
Map.keys(nested_api_data)
|> AsyncProcessor.execute(
fn branch_key ->
lower_level_api_data = Map.get(nested_api_data, branch_key)
updated_keys_list = update_keys_list(root_key, [branch_key])
extract_nested_keys(lower_level_api_data, updated_keys_list)
|> List.flatten()
end,
:flat_map
)
end
def extract_nested_keys(nested_api_data, key) when not is_map(nested_api_data), do: key
def extract_keys(api_response_json, target_data_key) when is_atom(target_data_key) do
{:ok, target_data} = extract_target_data(api_response_json, target_data_key)
{:ok, decoded_api_response_json} = Jason.decode(target_data)
Map.keys(decoded_api_response_json)
|> AsyncProcessor.execute(
fn key ->
nested_keys_list =
Map.get(decoded_api_response_json, key)
|> extract_nested_keys(key)
update_keys_list(nested_keys_list, [])
end,
:flat_map
)
|> Enum.uniq()
|> AsyncProcessor.execute(&String.to_atom(&1))
end
def decypher(raw_json, target_data, api_url, target_data_key) do
try do
Jason.decode(target_data, keys: :atoms!)
rescue
_ ->
LogBook.main("Lack of atomized versions of API's keys for #{api_url} prevents the safe conversion of this JSON. Creating missing key versions now . . . ", __MODULE__, 49, :warning)
extract_keys(raw_json, target_data_key)
LogBook.main(
"JSON key atomization complete. Getting fresh data from #{api_url} . . .",
__MODULE__,
55,
:warning
)
main(api_url, target_data_key)
end
end
def extract_target_data(api_response_json, target_data_key) when is_atom(target_data_key) do
target_data =
case target_data_key do
:none -> api_response_json
_ -> Map.get(api_response_json, target_data_key)
end
{:ok, target_data}
end
def update_log(api_url, error_message, {sleep_time, interval}, line_no) do
LogBook.main(
"Error: #{inspect(error_message)} - Retrying #{api_url} in #{sleep_time} #{interval} . . . ",
__MODULE__,
line_no,
:warning
)
end
def retry_api_url(api_url, error_message, attempt_count) do
cond do
attempt_count == 1 ->
update_log(api_url, error_message, {200, "ms"}, 14)
Process.sleep(200)
main(api_url, attempt_count)
attempt_count == 2 ->
update_log(api_url, error_message, {500, "ms"}, 20)
Process.sleep(500)
main(api_url, attempt_count)
attempt_count <= 29 ->
update_log(api_url, error_message, {1.5, "seconds"}, 26)
Process.sleep(1_500)
main(api_url, attempt_count)
attempt_count > 29 ->
raise "Error after #{attempt_count} failed attempts: #{inspect(error_message)}"
end
end
def process_api_response(
{:error, %HTTPoison.Error{id: nil, reason: _} = error_message},
api_url,
attempt_count
),
do: retry_api_url(api_url, error_message, attempt_count + 1)
def process_api_response({:ok, response}, _, _), do: {:ok, response}
def connect_to_api(api_url, attempt_count \\ 0) do
response = HTTPoison.get(api_url)
process_api_response(response, api_url, attempt_count)
end
def fetch_data(api_url) do
case connect_to_api(api_url) do
{:ok, _} = good_response ->
good_response
glitch ->
raise "#{__MODULE__}: Mishandled value: #{inspect(glitch)}"
end
end
def main(api_url, target_data_key) do
with {:ok, api_response_json} <- fetch_data(api_url),
{:ok, target_data} <- extract_target_data(api_response_json, target_data_key) do
decypher(api_response_json, target_data, api_url, target_data_key)
end
end
end
Any particular advice on how to increase readability, efficiency, execution speed, conformance to best practices, etc. would be greatly appreciated.
Trending in Questions
I having some trouble figuring out if I have set myself too strict of standards for my production server. Currently I can handle 75% of r...
New
Documentation
While reading the Scoped Routes section, I noticed that the documentation currently refers to a problem without explainin...
New
Hello,
I’m trying to build a basic Phoenix web-app, and I’d like to use Tailwind.
However, when I launch mix phx.server, I get an error...
New
Hi everyone,
I am toying with the idea of building a “match maker” for giving personal help to people that wants to start coding.
I sta...
New
I recently noticed that Elixir’s Logger defaults its primary log level to :debug when no :logger, :level application configuration is pre...
New
I’m working on a small exercise involving update_in/3, and I came up with this solution:
data = %{
name: "Periodic Table",
category:...
New
I’ve got trouble wrapping my head around the order in which functions are called in this snippet (from Phoenix’s authentication):
toke...
New
Other Trending Topics
Edit: 2026 May 15 - This post is archived.
Mob is alive!!
Main docs: mob v0.7.11 — Documentation
A bit of explanation for the slightly c...
New
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
I am happy to introduce the very α version of the new programming language compiled to BEAM.
Welcome Cure.
It has literally three kille...
New
Hobbes is a low-level distributed database for the Elixir programming language.
Hobbes provides a simple, safe, and scalable storage lay...
New
Hi everyone!
The first release candidate for the Expert language server project is now available!
We’ve published a press release detai...
New
A little off-topic, but I feel like people here have a good head on their shoulders.
I used to be quite good at making software. Was luc...
New
Categories:
Sub Categories:
Forums
Popular Tags
- #ecto
- #liveview
- #troubleshooting
- #learning-elixir
- #library
- #deployment
- #erlang
- #testing
- #genserver
- #mix
- #absinthe
- #remote-other
- #otp
- #plug
- #how-to-question
- #macros
- #postgres
- #elixirconf
- #channels
- #exunit
- #discussion
- #code-sync
- #podcasts
- #javascript
- #onsite
- #dialyzer
- #docker
- #authentication
- #umbrella
- #full-time-contract
- #podcasts-by-brainlid
- #ai
- #ecto-query
- #elixirconf-us
- #blog-post
- #elixir-ls
- #phoenix_html
- #iex
- #graphql
- #genstage
- #websockets
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #elixirconf-eu
- #api
- #forms
- #metaprogramming
- #hex











Showing Posts 1 to 8- Show Best Posts
- Show All (oldest first)
- Show All (newest first)
thiagomajesk
Hey @Maxximiliann! This is a really broad request with very little context, could you elaborate further?
It would be helpful to give a little more context to your code snippet, like what purpose does this code serve and what are you trying to accomplish with it… It would also be very nice if you could describe what have tried so far and if there are any particular parts you are seeking to improve.
This is the real combo, isn’t it
!? Jokes aside, It’s very hard to answer an open-ended question like that. So if possible, it would be better to be a little bit more specific.
al2o3cr
Functionality observations:
:flat_maphead will crash if any task returns{:ok, []}.:flat_maphead returns results reverse order compared to the othersexecutereturn a list, others return a yet-to-be-evaluatedStreamIMO this design adds a lot of complexity to try to cover every use case in one function. For instance, what would the
@specforexecute’slambda’s possible types look like?Consider splitting
executeinto multiple functions that each expect a callback with a specific shape.write_to_logwill always return{:ok, _}, so this entirewithstatement is useless.There is significant overhead involved in
Task.async_stream, especially compared to literally fetching a single atom from the atom table. This is not an efficient way to do things.Furthermore, the only things you should be doing in
Task.async_streamwith@max_concurrency System.schedulers_online() * 30are wait-heavy IO operation, since there are so many more processes than available CPUs.This function wears two hats at once: it formats the data with
inspectand it logs it. Consider splitting those responsibilities.IMO functions that dispatch on an argument (like
formatabove, orwrite_to_console) are fine, but it becomes a design smell when callers of those functions are also special-casing those arguments. Refactoring can tidy this up, for instance by aligning the structures and keeping the behavior to a single function. For instance:Maxximiliann
First of all, I want to thank you so much for all your astute observations; this is exactly the kind of feedback I was searching for. I’ve learned a lot from them!
Please elaborate on what you mean exactly by “overhead,” if you kindly would.
Thank you for this insight. How exactly can I tell whether an operation is IO wait-heavy or not?
dimitarvp
For the record I can’t exactly agree with @al2o3cr here because parallelizing work (one that’s inherently parallel, of course) still yields huge wins compared to trying to do single-thread async switching. I am not aware of the significant overhead he’s talking about but on a higher level I never regretted parallelizing things with Elixir. Of course some – perhaps not obvious for everyone – limitations must apply e.g. if you have some mere 10_000 items to process where each item does not interact with network or disk then doing this in parallel is highly likely to be not worth it. Which brings me to…
Modern CPUs spend most of their lives waiting on network or disk. If your parallel tasks involve I/O (so network or disks or any slow periphery) then they are a perfect candidate for high parallelization (meaning parallel task amount that’s bigger than your CPU threads).
…Although that has limitations as well e.g. it’s not very useful to spawn 100 tasks that all read/write from/to the same slow hard drive. As always, you should measure and take these factors into account.
@al2o3cr’s remark is valid as long as your tasks are actually CPU-bound; if they are then having more than CPU threads parallel tasks (
System.schedulers_online()is usually the same number) will not achieve any performance improvement. He’s right about that.The general rule of parallelization is:
Maxximiliann
Thanks so much for your keen insights. Could you kindly elaborate on how one goes about measuring for this?
cjbottaro
That’s going to set
@max_concurrencyto the number of processors of the computer that compiled the code, not the number of the processors of the computer that the application is run on.See compile time vs runtime config.
al2o3cr
“Overhead” as in “starting Tasks and awaiting the results isn’t free”. A quick benchmark:
gives results like:
Using
async_streamis 25x slower thanmap, and it gets EVEN SLOWER (70xmap!) when trying to use too much concurrency. This is particularly notable here becausemap_fundoes very very little, so all of this is overhead.Slowing
map_fundown by adding aProcess.sleep(2)changes the situation:Now every invocation of
map_funtakes a lot more “wall-clock time” than before, but approximately the same amount of “running on the CPU” time since it’s mostly asleep.Because of that,
async_streamis faster (by about a factor ofSystem.schedulers_online) - andtoo_many_async_streamis even faster, though with diminishing returns aftermax_concurrencyof about 64 or so.dimitarvp
Yep, parallelization should only be utilized when the tasks are either not trivial (take some CPU time) or just involve some waiting (I/O). In all other scenarios and with most list sizes you’re better off just doing
Enum.mapover the list.