Class: Sidekiq::ElasticsearchProcessQueue

Inherits:
Object
  • Object
show all
Includes:
Worker
Defined in:
app/workers/sidekiq/elasticsearch_process_queue.rb

Instance Method Summary collapse

Instance Method Details

#on_success(_status, options) ⇒ Object



31
32
33
34
35
36
37
38
# File 'app/workers/sidekiq/elasticsearch_process_queue.rb', line 31

def on_success(_status, options)
  index_class_name = options["index_class_name"]
  record_ids_count = options["record_ids_count"]
  key = options["key"]
  starting_time = options["starting_time"]
  ::Stat.save_to_stats(Stat.new, "Sidekiq::ElasticsearchProcessQueue", "elasticsearch", index_class_name, { record_ids_count: record_ids_count }, starting_time, ::Process.clock_gettime(::Process::CLOCK_MONOTONIC)) if ::HasHelpers::Feature.active?(:elasticsearch_reindex_stat)
  ::HasHelpers::Index.after_index(key)
end

#perform(index_class_name, index_name) ⇒ Object



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