Current section
Files
Jump to
Current section
Files
lib/analytics_3i/workers/collector.ex
defmodule Workers.Collector do
@moduledoc false
@size 15
@defaults %{default_options: %{slop: 50, size: @size, withTitle: true, withRank: true}}
use GenServer
require Logger
def start_link do
GenServer.start_link(__MODULE__, [], name: :collector_server, restart: :temporary)
end
def init(_opts) do
tasks = :ets.new(__MODULE__, [:set, :named_table])
{:ok, tasks}
end
def lookup, do: GenServer.call(:collector_server, {:lookup})
def lookup(id), do: GenServer.call(:collector_server, {:lookup, id})
def lookup_ref(tasks, id) do
case :ets.lookup(tasks, id) do
[{^id, pid, otrasl, tech, queries, status, count}] ->
{:ok, %{ pid: pid, otrasl: otrasl, tech: tech, queries: queries, status: status, id: id, count: count}}
[] -> :error
end
end
def last_task_get(tasks), do: {:ok, :ets.last(tasks)}
def last_task, do: GenServer.call(:collector_server, {:last_task})
def delete_task(task) do
GenServer.cast(:collector_server, {:delete_task, task})
end
@doc false
def new(requests) do
GenServer.call(:collector_server, {:new, requests})
end
@doc false
def execute do
GenServer.cast(:collector_server, {:execute})
end
def handle_call({:lookup}, _, tasks), do: {:reply, self(), tasks}
def handle_call({:lookup, id}, _, tasks), do: {:reply, lookup_ref(tasks, id), tasks}
def handle_call({:last_task}, _, tasks), do: {:reply, last_task_get(tasks), tasks}
def handle_call({:new, requests}, {from, _}, tasks) do
requests |> Enum.each(fn {otrasl, techs} ->
techs |> Enum.each(fn {tech, lists} ->
%{ status: :init, otrasl: otrasl, tech: tech, queries: lists, id: UUID.uuid1(), pid: from, count: 0}
|> insert_task(tasks)
end)
end)
{:reply, :ok, tasks}
end
def handle_cast({:delete_task, task}, tasks) do
true = :ets.delete(tasks, task[:id])
{:noreply, tasks}
end
def handle_cast({:execute}, tasks) do
stack_size(tasks)
{:ok, task_id} = last_task_get(tasks)
case task_id do
"" <> id ->
case lookup_ref(tasks, id) do
{:ok, task} ->
Logger.debug "START #{id} #{task[:otrasl]} #{task[:tech]}"
result = iterate_tech!(task[:queries])
write_tech_csv!(task[:otrasl], task[:tech], result)
task = %{ task | count: Enum.count(result), status: :success }
Logger.debug "FINISH #{id} #{task[:otrasl]} #{task[:tech]}"
Workers.Collector.delete_task(task)
finish(task)
:error ->
:noop
end
_ -> :noop
end
{:noreply, tasks}
end
def stack_size(tasks) do
info = :ets.info(tasks)
Logger.debug "STACK SIZE: #{info[:size]}"
end
defp write_tech_csv!(otrasl, tech, values) do
Logger.debug "WRITING START #{otrasl}, #{tech}"
File.open("output/3i_#{otrasl}_#{tech}.csv", [:write, :utf8], fn(file) ->
values
|> CSV.encode(headers: [:id, :created_date, :title, :url, :rank, :language, :provider])
|> Enum.each(&IO.write(file, &1))
end)
Logger.debug "WRITING FINISH #{otrasl}, #{tech}"
end
defp iterate_tech!(lists) do
lists
|> Enum.reduce([], fn {num, list}, data ->
Logger.debug "NUM START #{num}"
queries = list["queries"]
sd = Enum.reduce(queries, [], fn pre_query, small_data ->
query = "(#{list["required"]}) AND (#{pre_query})"
small_data ++ collect!(query)
end)
Logger.debug "NUM FINISH #{num}"
data ++ sd
end)
|> DomainFilter.log
|> DomainFilter.filter
|> DuplicatesFilter.log
|> DuplicatesFilter.filter
# |> JaroWinklerFilter.log
# |> JaroWinklerFilter.filter
end
defp collect!(query, old_data \\ [], page \\ 1) do
Process.sleep(1500)
{:ok, data} = Analytics3i.search(query, Map.put(@defaults, :default_options, Map.merge(@defaults[:default_options], %{ from: @size * (page - 1) })))
if (data[:total] || 0) > (page * @size), do: collect!(query, old_data ++ data[:articles], page + 1), else: old_data ++ data[:articles]
end
defp finish(%{pid: pid} = task) do
send(pid, {:collector_finish, Map.delete(task, :queries)})
end
defp insert_task(%{id: id, pid: pid} = prepare_task, tasks) do
true = :ets.insert(tasks, {id, pid, prepare_task[:otrasl], prepare_task[:tech], prepare_task[:queries], prepare_task[:status], prepare_task[:count]})
end
end