8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
|
# File 'app/workers/sidekiq/elasticsearch_process_queue.rb', line 8
def perform(index_class_name, index_name)
starting_time = ::Process.clock_gettime(::Process::CLOCK_MONOTONIC)
index_klass = index_class_name.constantize
limit = 100
is_keyword_search = index_name.include?("keyword")
reindex_queue = is_keyword_search ? index_klass.keyword_reindex_queue : index_klass.reindex_queue
records = reindex_queue.reserve(limit: limit)
if records.any?
batch = Sidekiq::Batch.new
record_ids = records.map { |record| record[:record_id] }
key = records.map { |record| record[:key] }.uniq
routing_map = records.each_with_object({}) { |record, h| h[record[:record_id]] = record[:routing_key] }
batch.on(:success, self.class, "index_class_name" => index_class_name, "record_ids_count" => record_ids.size, "key" => key, "starting_time" => starting_time)
batch.jobs do
::Sidekiq::ElasticsearchProcessBatch.perform_async(index_klass.name, index_name, record_ids, routing_map) if ::Searchkick.client.indices.exists_alias?(name: index_name)
end
if record_ids.size == limit
::Sidekiq::ElasticsearchProcessQueue.perform_async(index_class_name, index_name)
end
end
end
|