Dowser. Elasticsearch. Streamer
(Dowser.Elasticsearch v0.4.1)
View Source
Turns a search into an Elixir Stream, for walking more documents than fit
in one response.
%{query: %{match_all: %{}}, size: 1_000}
|> Dowser.Elasticsearch.Streamer.stream(index: "posts")
|> Stream.map(& &1["_source"])
|> Enum.each(&process/1)Each element is a hit — _index, _id, _source and the rest — cast
exactly as Dowser.Elasticsearch.Search.search/2 would cast it.
How it walks
A point in time
pins the index against concurrent writes, and search_after pages through
it. _shard_doc is appended to the sort as a tiebreaker, which is what makes
the paging deterministic; a sort of your own is kept and sorted on first.
search_after is sequential by construction — a page's cursor is the last
hit of the page before it — so one stream cannot fetch pages in parallel.
stream_slices/4 is how you get parallelism; see below.
Options
:index— the index to open the point in time on. Required unless:pitis given. It is not sent with the searches themselves: Elasticsearch rejects a search that names both an index and a PIT.:pit— an existing point-in-time id to walk instead of opening one. It is checked once before the walk starts, and left open afterwards — whoever opened it closes it. Without it, the stream opens its own and closes it when enumeration ends, however it ends.:keep_alive— how long Elasticsearch holds the point in time, extended on every search. Defaults to"1m". Long enough to cover the gap between two pages, not the whole walk.:verify_pit— whether a:pityou passed is checked before the walk starts.trueby default;falsewhen you just opened it and know it is alive.
Everything else is forwarded to Dowser.Elasticsearch.Search, so :context,
:codec, :keys and :http_opts all work as usual.
A query with no size gets 1000, not Elasticsearch's default of 10 —
at 10 hits per round trip a million documents is a hundred thousand
requests. Set your own when you have a reason to.
Slicing
stream_slices/4 splits the point in time into disjoint subsets and
walks them at once. Elasticsearch divides first across shards, then within
each shard by contiguous ranges of Lucene document ids, so the natural
ceiling is your shard count — more slices than shards subdivides a shard
rather than adding parallelism.
%{query: %{match_all: %{}}, size: 1_000}
|> Dowser.Elasticsearch.Streamer.stream_slices(4, &Enum.count/1, index: "posts")
|> Enum.sum()Note what it takes: a function, not a stream. Each slice is consumed inside its own task, because a lazy stream handed back out of a task would run every page in the caller — concurrent in name only.
To spread a walk across nodes instead, open the point in time yourself and
give each node one slice. slice is an ordinary search body field, so it
goes in the query rather than the options:
stream(%{query: ..., slice: %{id: 2, max: 8}}, pit: pit_id)Failure
There is no stream!/2. Every other function in this package comes in a pair
because it returns {:ok, result} or {:error, error} and the bang variant
unwraps it; a stream has nothing to unwrap, and nothing has happened yet when
it is built. Errors surface on enumeration, by raising — which is what the
bang variant would have done anyway.
The stream raises on the first failed request, and closes a point in time it
opened on the way out — on normal completion, on Enum.take/2, and on an
exception. It cannot close one if the enumerating process is killed outright;
:keep_alive is the backstop there.
Summary
Functions
Streams every hit a search matches.
Walks slice_nbr slices of one point in time at once, running stream_fn
over each.
Types
@type query() :: map()
Functions
@spec stream(query(), keyword()) :: Enumerable.t()
Streams every hit a search matches.
See the module documentation for the options. The stream is lazy: nothing is requested, and no point in time is opened, until it is enumerated.
@spec stream_slices(query(), pos_integer(), (Enumerable.t() -> term()), keyword()) :: Enumerable.t()
Walks slice_nbr slices of one point in time at once, running stream_fn
over each.
Both are positional because both are required: there is no sensible default for how many slices to open — the useful number is your shard count, which this cannot know — and a function is what the whole call is for.
Note what it yields: one result per slice, not a stream of hits like
stream/2. The hits are stream_fn's to consume.
stream_fn receives a slice's stream and is called inside the task that
owns it, so the hits never cross a process boundary — which is the whole
point: returning a lazy stream from a task would build it there and then run
every page back in the caller.
%{query: %{match_all: %{}}, size: 1_000}
|> Dowser.Elasticsearch.Streamer.stream_slices(4, &Enum.count/1, index: "posts")
|> Enum.sum()The result is a stream of whatever stream_fn returned, one per slice, so
keep those small — a count, a sum, :ok. A slice that returns everything it
read defeats the streaming.
A slice already in the query body is the base every slice is built on, so
anything else the split needs travels with it; id and max are computed
here and win:
stream_slices(%{query: ..., slice: %{field: "_id"}}, 7, &f/1, index: "posts")
# each walk carries %{"field" => "_id", "id" => 0..6, "max" => 7}All slices share one point in time, opened here and closed when the stream
finishes, however it finishes. Pass :pit to use one of your own; it is
checked once, and left open.
Options
Takes everything stream/2 does, plus:
:max_concurrency— how many slices run at once. Defaults toslice_nbr, so every slice you asked for is actually in flight.Note this is not
Task.async_stream/3's default ofSystem.schedulers_online/0, nor capped by it. A slice spends its time waiting on Elasticsearch, and a process blocked on a socket occupies no scheduler — so the limit that matters is what the cluster will take, not how many cores this machine has. Capping 32 slices at 10 schedulers makes the walk four times slower for nothing.Lower it when
stream_fnis the expensive part rather than the fetching, or to be gentler on the cluster.:timeout— per slice, not per request. Defaults to:infinity, since a slice runs as long as it takes to walk;Task.async_stream/3would otherwise give up after five seconds.:ordered—falseby default, so a finished slice is not held back by a slower one.
A slice that fails raises: an exception from its own walk, or a
RuntimeError naming the slice if the task exited or timed out.