This commit is contained in:
@@ -11,7 +11,9 @@ defmodule LoLAnalytics.Application do
|
||||
LoLAnalytics.Repo,
|
||||
{DNSCluster, query: Application.get_env(:lol_analytics, :dns_cluster_query) || :ignore},
|
||||
{Phoenix.PubSub, name: LoLAnalytics.PubSub},
|
||||
{Task.Supervisor, name: LoLAnalytics.TaskSupervisor}
|
||||
{Task.Supervisor, name: LoLAnalytics.TaskSupervisor},
|
||||
{LolAnalytics.MatchProcessor.MatchesBroadwayProcessor, []},
|
||||
{LolAnalytics.MatchProcessor.MatchesProducer, []}
|
||||
# Start a worker by calling: LoLAnalytics.Worker.start_link(arg)
|
||||
# {LoLAnalytics.Worker, arg}
|
||||
]
|
||||
|
||||
@@ -61,10 +61,25 @@ defmodule LolAnalytics.Dimensions.Match.MatchRepo do
|
||||
Repo.all(MatchSchema)
|
||||
end
|
||||
|
||||
def list_unprocessed_matches(limit, queue \\ 420) do
|
||||
query =
|
||||
from m in MatchSchema,
|
||||
where:
|
||||
(m.fact_champion_picked_item_status == 0 or
|
||||
m.fact_champion_picked_summoner_spell_status == 0 or
|
||||
m.fact_champion_played_game_status == 0) and
|
||||
m.queue_id == ^queue,
|
||||
order_by: [desc: m.updated_at],
|
||||
limit: ^limit
|
||||
|
||||
Repo.all(query)
|
||||
end
|
||||
|
||||
@type process_status :: :not_processed | :processed | :error
|
||||
defp process_status_atom_to_db(:not_processed), do: 0
|
||||
defp process_status_atom_to_db(:enqueued), do: 1
|
||||
defp process_status_atom_to_db(:processed), do: 2
|
||||
defp process_status_atom_to_db(:error), do: 3
|
||||
defp process_status_atom_to_db(:error_match_not_found), do: 4
|
||||
defp process_status_atom_to_db(_), do: raise("Invalid processing status")
|
||||
end
|
||||
|
||||
@@ -1,8 +1,9 @@
|
||||
defmodule LolAnalytics.Facts.ChampionPickedItem.FactProcessor do
|
||||
require Logger
|
||||
|
||||
alias LolAnalytics.Dimensions.Match.MatchSchema
|
||||
alias LolAnalytics.Dimensions.Match.MatchRepo
|
||||
alias LolAnalytics.Facts.ChampionPickedItem.Repo
|
||||
require Logger
|
||||
@behaviour LolAnalytics.Facts.FactBehaviour
|
||||
|
||||
@doc """
|
||||
|
||||
@@ -21,6 +22,24 @@ defmodule LolAnalytics.Facts.ChampionPickedItem.FactProcessor do
|
||||
end
|
||||
end
|
||||
|
||||
@spec process_match(%MatchSchema{}) :: :ok | {:error, String.t()}
|
||||
def process_match(match) do
|
||||
match_url = "http://192.168.1.55:9000/ranked/#{match.patch_number}/#{match.match_id}.json"
|
||||
|
||||
with {:ok, %HTTPoison.Response{status_code: 200, body: body}} <-
|
||||
HTTPoison.get(match_url),
|
||||
{:ok, decoded_match} <- Poison.decode(body, as: %LoLAPI.Model.MatchResponse{}) do
|
||||
process_game_data(decoded_match)
|
||||
MatchRepo.update(match, %{fact_champion_picked_item_status: :processed})
|
||||
:ok
|
||||
else
|
||||
_ ->
|
||||
MatchRepo.update(match, fact_champion_picked_item_status: :error_match_not_found)
|
||||
Logger.error("Could not process data from #{match_url} for ChampionPickedItem")
|
||||
{:error, "Could not process data from #{match_url}"}
|
||||
end
|
||||
end
|
||||
|
||||
defp process_game_data(decoded_match) do
|
||||
participants = decoded_match.info.participants
|
||||
version = extract_game_version(decoded_match)
|
||||
@@ -57,8 +76,6 @@ defmodule LolAnalytics.Facts.ChampionPickedItem.FactProcessor do
|
||||
end)
|
||||
end
|
||||
end)
|
||||
|
||||
MatchRepo.update(match, %{fact_champion_picked_item_status: :processed})
|
||||
end
|
||||
|
||||
defp extract_game_version(game_data) do
|
||||
|
||||
+11
-10
@@ -1,22 +1,25 @@
|
||||
defmodule LolAnalytics.Facts.ChampionPickedSummonerSpell.FactProcessor do
|
||||
@behaviour LolAnalytics.Facts.FactBehaviour
|
||||
|
||||
require Logger
|
||||
|
||||
alias LolAnalytics.Dimensions.Match.MatchSchema
|
||||
alias LolAnalytics.Dimensions.Match.MatchRepo
|
||||
alias LolAnalytics.Facts.ChampionPickedSummonerSpell
|
||||
|
||||
@impl true
|
||||
@spec process_game_at_url(String.t()) :: any()
|
||||
def process_game_at_url(url) do
|
||||
@spec process_match(%MatchSchema{}) :: :ok | {:error, String.t()}
|
||||
def process_match(match) do
|
||||
match_url = "http://192.168.1.55:9000/ranked/#{match.patch_number}/#{match.match_id}.json"
|
||||
|
||||
with {:ok, %HTTPoison.Response{status_code: 200, body: body}} <-
|
||||
HTTPoison.get(url),
|
||||
HTTPoison.get(match_url),
|
||||
{:ok, decoded_match} <- Poison.decode(body, as: %LoLAPI.Model.MatchResponse{}) do
|
||||
process_game_data(decoded_match)
|
||||
MatchRepo.update(match, %{fact_champion_picked_summoner_spell_status: :processed})
|
||||
:ok
|
||||
else
|
||||
_ ->
|
||||
Logger.error("Could not process data from #{url} for ChampionPickedSummonerSpell")
|
||||
{:error, "Could not process data from #{url}"}
|
||||
MatchRepo.update(match, fact_champion_picked_summoner_spell_status: :error_match_not_found)
|
||||
Logger.error("Could not process data from #{match_url} for ChampionPickedItem")
|
||||
{:error, "Could not process data from #{match_url}"}
|
||||
end
|
||||
end
|
||||
|
||||
@@ -64,8 +67,6 @@ defmodule LolAnalytics.Facts.ChampionPickedSummonerSpell.FactProcessor do
|
||||
ChampionPickedSummonerSpell.Repo.insert(attrs_spell_2)
|
||||
end
|
||||
end)
|
||||
|
||||
MatchRepo.update(match, %{fact_champion_picked_summoner_spell_status: :processed})
|
||||
end
|
||||
|
||||
defp extract_game_version(game_data) do
|
||||
|
||||
@@ -1,20 +1,24 @@
|
||||
defmodule LolAnalytics.Facts.ChampionPlayedGame.FactProcessor do
|
||||
alias LolAnalytics.Dimensions.Match.MatchRepo
|
||||
require Logger
|
||||
|
||||
@behaviour LolAnalytics.Facts.FactBehaviour
|
||||
alias LolAnalytics.Dimensions.Match.MatchSchema
|
||||
alias LolAnalytics.Dimensions.Match.MatchRepo
|
||||
|
||||
@spec process_match(%MatchSchema{}) :: :ok | {:error, String.t()}
|
||||
def process_match(match) do
|
||||
match_url = "http://192.168.1.55:9000/ranked/#{match.patch_number}/#{match.match_id}.json"
|
||||
|
||||
@impl true
|
||||
@spec process_game_at_url(String.t()) :: none()
|
||||
def process_game_at_url(url) do
|
||||
with {:ok, %HTTPoison.Response{status_code: 200, body: body}} <-
|
||||
HTTPoison.get(url),
|
||||
HTTPoison.get(match_url),
|
||||
{:ok, decoded_match} <- Poison.decode(body, as: %LoLAPI.Model.MatchResponse{}) do
|
||||
process_game_data(decoded_match)
|
||||
MatchRepo.update(match, %{fact_champion_played_game_status: :processed})
|
||||
:ok
|
||||
else
|
||||
_ ->
|
||||
Logger.error("Could not process data from #{url} for ChampionPlayedGame")
|
||||
{:error, "Could not process data from #{url}"}
|
||||
MatchRepo.update(match, fact_champion_played_game_status: :error_match_not_found)
|
||||
Logger.error("Could not process data from #{match_url} for ChampionPickedItem")
|
||||
{:error, "Could not process data from #{match_url}"}
|
||||
end
|
||||
end
|
||||
|
||||
@@ -48,8 +52,6 @@ defmodule LolAnalytics.Facts.ChampionPlayedGame.FactProcessor do
|
||||
LolAnalytics.Facts.ChampionPlayedGame.Repo.insert(attrs)
|
||||
end
|
||||
end)
|
||||
|
||||
MatchRepo.update(match, %{fact_champion_played_game_status: :processed})
|
||||
end
|
||||
|
||||
defp extract_game_version(game_data) do
|
||||
|
||||
@@ -1,3 +0,0 @@
|
||||
defmodule LolAnalytics.Facts.FactBehaviour do
|
||||
@callback process_game_at_url(String.t()) :: any()
|
||||
end
|
||||
@@ -1,34 +1,18 @@
|
||||
defmodule LolAnalytics.Facts.FactsRunner do
|
||||
alias LolAnalytics.Facts
|
||||
|
||||
def analyze_by_patch(patch) do
|
||||
Storage.MatchStorage.S3MatchStorage.stream_files("ranked", patch: patch)
|
||||
|> peach(fn %{key: path} ->
|
||||
get_facts()
|
||||
|> Enum.each(fn fact_runner ->
|
||||
apply(fact_runner, ["http://192.168.1.55:9000/ranked/#{path}"])
|
||||
end)
|
||||
def analyze_match(match) do
|
||||
get_facts()
|
||||
|> Enum.each(fn fact_runner ->
|
||||
apply(fact_runner, [match])
|
||||
end)
|
||||
end
|
||||
|
||||
def analyze_all_matches do
|
||||
Storage.MatchStorage.S3MatchStorage.stream_files("ranked")
|
||||
|> peach(fn %{key: path} ->
|
||||
get_facts()
|
||||
|> Enum.each(fn fact_runner ->
|
||||
apply(fact_runner, ["http://192.168.1.55:9000/ranked/#{path}"])
|
||||
end)
|
||||
end)
|
||||
end
|
||||
|
||||
def analyze_match() do
|
||||
end
|
||||
|
||||
def get_facts() do
|
||||
[
|
||||
&Facts.ChampionPickedSummonerSpell.FactProcessor.process_game_at_url/1,
|
||||
&Facts.ChampionPlayedGame.FactProcessor.process_game_at_url/1,
|
||||
&Facts.ChampionPickedItem.FactProcessor.process_game_at_url/1
|
||||
&Facts.ChampionPickedSummonerSpell.FactProcessor.process_match/1,
|
||||
&Facts.ChampionPlayedGame.FactProcessor.process_match/1,
|
||||
&Facts.ChampionPickedItem.FactProcessor.process_match/1
|
||||
]
|
||||
end
|
||||
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
defmodule LolAnalytics.MatchProcessor.MatchesBroadwayProcessor do
|
||||
alias LolAnalytics.Facts.FactsRunner
|
||||
use Broadway
|
||||
|
||||
def start_link(opts) do
|
||||
Broadway.start_link(__MODULE__,
|
||||
name: __MODULE__,
|
||||
processors: [default: []],
|
||||
producer: [
|
||||
module: {LolAnalytics.MatchProcessor.MatchesProducer, []},
|
||||
rate_limiting: [
|
||||
interval: 1000,
|
||||
allowed_messages: 40
|
||||
]
|
||||
]
|
||||
)
|
||||
end
|
||||
|
||||
@impl Broadway
|
||||
def handle_message(_processor, message, _context) do
|
||||
message.data
|
||||
# build_match_url(message.data.queue_id, message.data.patch_number, message.data.match_id)
|
||||
|> FactsRunner.analyze_match()
|
||||
|
||||
message
|
||||
end
|
||||
|
||||
defp build_match_url(queue, patch_id, match_id) do
|
||||
"http://192.168.1.55:9000/#{queue_to_dir(queue)}/#{patch_id}/#{match_id}.json"
|
||||
end
|
||||
|
||||
defp queue_to_dir(420), do: "ranked"
|
||||
defp queue_to_dir(_), do: "ranked"
|
||||
end
|
||||
@@ -0,0 +1,33 @@
|
||||
defmodule LolAnalytics.MatchProcessor.MatchesProducer do
|
||||
use GenStage
|
||||
|
||||
@impl GenStage
|
||||
def init(opts) do
|
||||
{:producer, opts}
|
||||
end
|
||||
|
||||
def start_link(opts) do
|
||||
GenStage.start_link(__MODULE__, :ok)
|
||||
end
|
||||
|
||||
@impl GenStage
|
||||
def handle_demand(demand, state) do
|
||||
matches = query_unprocessed_matches(demand)
|
||||
|
||||
{:noreply, matches, state}
|
||||
end
|
||||
|
||||
defp query_unprocessed_matches(demand) when demand <= 0, do: []
|
||||
|
||||
defp query_unprocessed_matches(demand) do
|
||||
LolAnalytics.Dimensions.Match.MatchRepo.list_unprocessed_matches(demand)
|
||||
|> Enum.map(&broadway_transform/1)
|
||||
end
|
||||
|
||||
defp broadway_transform(match) do
|
||||
%Broadway.Message{
|
||||
data: match,
|
||||
acknowledger: Broadway.NoopAcknowledger.init()
|
||||
}
|
||||
end
|
||||
end
|
||||
@@ -1,26 +0,0 @@
|
||||
defmodule LolAnalytics.MatchesProcessor do
|
||||
use GenServer
|
||||
|
||||
def init(init_args) do
|
||||
{:ok, init_args}
|
||||
end
|
||||
|
||||
@doc """
|
||||
iex> LolAnalytics.MatchesProcessor.process_for_patch "14.12.593.5894"
|
||||
"""
|
||||
def process_for_patch(patch) do
|
||||
Task.Supervisor.async(LoLAnalytics.TaskSupervisor, fn ->
|
||||
LolAnalytics.Facts.FactsRunner.analyze_by_patch(patch)
|
||||
end)
|
||||
end
|
||||
|
||||
def process_all_matches() do
|
||||
Task.Supervisor.async(LoLAnalytics.TaskSupervisor, fn ->
|
||||
LolAnalytics.Facts.FactsRunner.analyze_all_matches()
|
||||
end)
|
||||
end
|
||||
|
||||
def get_running_processes() do
|
||||
Task.Supervisor.children(LoLAnalytics.TaskSupervisor)
|
||||
end
|
||||
end
|
||||
Reference in New Issue
Block a user