forked from AkkomaGang/akkoma
Make the indexing batch differently and more, show number indexed
This commit is contained in:
parent
e5ac2ffa07
commit
abf82a63ec
1 changed files with 39 additions and 26 deletions
|
@ -28,33 +28,46 @@ def run(["index"]) do
|
||||||
])
|
])
|
||||||
)
|
)
|
||||||
|
|
||||||
Pleroma.Repo.chunk_stream(
|
chunk_size = 100_000
|
||||||
from(Pleroma.Object,
|
|
||||||
# Only index public posts which are notes and have some text
|
|
||||||
where:
|
|
||||||
fragment("data->>'type' = 'Note'") and
|
|
||||||
fragment("LENGTH(data->>'source') > 0") and
|
|
||||||
fragment("data->'to' \\? ?", ^Pleroma.Constants.as_public())
|
|
||||||
),
|
|
||||||
200,
|
|
||||||
:batches
|
|
||||||
)
|
|
||||||
|> Stream.map(fn objects ->
|
|
||||||
Enum.map(objects, fn object ->
|
|
||||||
data = object.data
|
|
||||||
%{id: object.id, source: data["source"], ap: data["id"]}
|
|
||||||
end)
|
|
||||||
end)
|
|
||||||
|> Stream.each(fn objects ->
|
|
||||||
{:ok, _} =
|
|
||||||
Pleroma.HTTP.post(
|
|
||||||
"#{endpoint}/indexes/objects/documents",
|
|
||||||
Jason.encode!(objects)
|
|
||||||
)
|
|
||||||
|
|
||||||
IO.puts("Indexed #{Enum.count(objects)} entries")
|
Pleroma.Repo.transaction(
|
||||||
end)
|
fn ->
|
||||||
|> Stream.run()
|
Pleroma.Repo.stream(
|
||||||
|
from(Pleroma.Object,
|
||||||
|
# Only index public posts which are notes and have some text
|
||||||
|
where:
|
||||||
|
fragment("data->>'type' = 'Note'") and
|
||||||
|
fragment("LENGTH(data->>'source') > 0") and
|
||||||
|
fragment("data->'to' \\? ?", ^Pleroma.Constants.as_public()),
|
||||||
|
order_by: fragment("data->'published' DESC")
|
||||||
|
),
|
||||||
|
timeout: :infinity
|
||||||
|
)
|
||||||
|
|> Stream.chunk_every(chunk_size)
|
||||||
|
|> Stream.transform(0, fn objects, acc ->
|
||||||
|
new_acc = acc + Enum.count(objects)
|
||||||
|
|
||||||
|
IO.puts("Indexed #{new_acc} entries")
|
||||||
|
|
||||||
|
{[objects], new_acc}
|
||||||
|
end)
|
||||||
|
|> Stream.map(fn objects ->
|
||||||
|
Enum.map(objects, fn object ->
|
||||||
|
data = object.data
|
||||||
|
%{id: object.id, source: data["source"], ap: data["id"]}
|
||||||
|
end)
|
||||||
|
end)
|
||||||
|
|> Stream.each(fn objects ->
|
||||||
|
{:ok, _} =
|
||||||
|
Pleroma.HTTP.post(
|
||||||
|
"#{endpoint}/indexes/objects/documents",
|
||||||
|
Jason.encode!(objects)
|
||||||
|
)
|
||||||
|
end)
|
||||||
|
|> Stream.run()
|
||||||
|
end,
|
||||||
|
timeout: :infinity
|
||||||
|
)
|
||||||
end
|
end
|
||||||
|
|
||||||
def run(["clear"]) do
|
def run(["clear"]) do
|
||||||
|
|
Loading…
Reference in a new issue