Module: Deimos::ActiveRecordConsume::BatchConsumption
- Includes:
- Consume::BatchConsumption
- Included in:
- Deimos::ActiveRecordConsumer
- Defined in:
- lib/deimos/active_record_consume/batch_consumption.rb,
sig/defs.rbs
Overview
Methods for consuming batches of messages and saving them to the database in bulk ActiveRecord operations.
Instance Method Summary collapse
-
#compact_messages(batch) ⇒ ::Array[Message]
Compact a batch of messages, taking only the last message for each unique key.
-
#consume_batch ⇒ void
Handle a batch of Kafka messages.
-
#deleted_query(records) ⇒ ActiveRecord::Relation
Create an ActiveRecord relation that matches all of the passed records.
-
#key_columns(_klass) ⇒ ::Array[String]
Get the set of attribute names that uniquely identify messages in the batch.
-
#record_key(key) ⇒ ::Hash[untyped, untyped]
Get unique key for the ActiveRecord instance from the incoming key.
-
#remove_records(messages) ⇒ void
Delete any records with a tombstone.
-
#uncompacted_update(messages) ⇒ void
Perform database operations for a batch of messages without compaction.
-
#update_database(messages) ⇒ void
Perform database operations for a group of messages.
-
#upsert_records(messages) ⇒ void
Upsert any non-deleted records records to either be updated or inserted.
Methods included from Consume::BatchConsumption
Instance Method Details
#compact_messages(batch) ⇒ ::Array[Message]
Compact a batch of messages, taking only the last message for each unique key.
@param batch — Batch of messages.
@return — Compacted batch.
126 127 128 129 130 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 126 def (batch) return batch unless batch.first&.key.present? batch.reverse.uniq(&:key).reverse! end |
#consume_batch ⇒ void
This method returns an undefined value.
Handle a batch of Kafka messages. Batches are split into "slices", which are groups of independent messages that can be processed together in a single database operation. If two messages in a batch have the same key, we cannot process them in the same operation as they would interfere with each other. Thus they are split
@param payloads — Decoded payloads
@param metadata — Information about batch, including keys.
28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 28 def consume_batch filtered = .select { || () } skipped_count = .size - filtered.size if skipped_count.positive? Deimos::Logging.log_debug( message: 'Skipping processing of messages in batch', skipped_count: skipped_count ) end = filtered.map { |p| Deimos::Message.new(p.payload, key: p.key) } tag = topic.name Deimos.config.tracer.active_span.set_tag('topic', tag) Karafka.monitor.instrument('deimos.ar_consumer.consume_batch', { topic: tag }) do failures = if @compacted && .map(&:key).compact.any? update_database(()) else uncompacted_update() end # Raised only once every slice and group has been attempted, so that a message which # can't be persisted never stops the rest of the batch from being saved. raise BatchFallbackError, failures if failures.any? end post_process_batch() end |
#deleted_query(records) ⇒ ActiveRecord::Relation
Create an ActiveRecord relation that matches all of the passed records. Used for bulk deletion.
@param records — List of messages.
@return — Matching relation.
97 98 99 100 101 102 103 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 97 def deleted_query(records) keys = records. map { |m| record_key(m.key)[@klass.primary_key] }. compact @klass.unscoped.where(@klass.primary_key => keys) end |
#key_columns(_klass) ⇒ ::Array[String]
Get the set of attribute names that uniquely identify messages in the batch. Requires at least one record.
@param records — Non-empty list of messages.
@return — List of attribute names.
65 66 67 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 65 def key_columns(_klass) nil end |
#record_key(key) ⇒ ::Hash[untyped, untyped]
Get unique key for the ActiveRecord instance from the incoming key. Override this method (with super) to customize the set of attributes that uniquely identifies each record in the database.
@param key — The encoded key.
@return — The key attributes.
81 82 83 84 85 86 87 88 89 90 91 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 81 def record_key(key) if key.nil? {} elsif key.is_a?(Hash) || key.is_a?(SchemaClass::Record) self.key_converter.convert(key) elsif self.topic.key_config[:field].nil? { @klass.primary_key => key } else { self.topic.key_config[:field].to_s => key } end end |
#remove_records(messages) ⇒ void
This method returns an undefined value.
Delete any records with a tombstone. deleted records.
@param messages — List of messages for a group of
338 339 340 341 342 343 344 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 338 def remove_records() Deimos::Utils::DeadlockRetry.wrap(Deimos.config.tracer.active_span.get_tag('topic')) do clause = deleted_query() clause.delete_all end end |
#uncompacted_update(messages) ⇒ void
This method returns an undefined value.
Perform database operations for a batch of messages without compaction. All messages are split into slices containing only unique keys, and each slice is handles as its own batch.
@param messages — List of messages.
138 139 140 141 142 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 138 def uncompacted_update() BatchSlicer. slice(). flat_map(&method(:update_database)) end |
#update_database(messages) ⇒ void
This method returns an undefined value.
Perform database operations for a group of messages. All messages with payloads are passed to upsert_records. All tombstones messages are passed to remove_records.
@param messages — List of messages.
150 151 152 153 154 155 156 157 158 159 160 161 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 150 def update_database() # Find all upserted records (i.e. that have a payload) and all # deleted record (no payload) removed, upserted = .partition { |m| delete_record?(m) } max_db_batch_size = self.class.config[:max_db_batch_size] upsert_groups = max_db_batch_size ? upserted.each_slice(max_db_batch_size).to_a : [upserted] remove_groups = max_db_batch_size ? removed.each_slice(max_db_batch_size).to_a : [removed] upsert_groups.reject(&:empty?).flat_map { |group| upsert_records(group) } + remove_groups.reject(&:empty?).flat_map { |group| remove_group(group) } end |
#upsert_records(messages) ⇒ void
This method returns an undefined value.
Upsert any non-deleted records records to either be updated or inserted.
@param messages — List of messages for a group of
170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 |
# File 'lib/deimos/active_record_consume/batch_consumption.rb', line 170 def upsert_records() record_list = build_records() invalid = filter_records(record_list) if invalid.any? Karafka.monitor.instrument('deimos.batch_consumption.invalid_records', { records: invalid, consumer: self.class }) end return [] if record_list.empty? key_col_proc = self.method(:key_columns).to_proc col_proc = self.method(:columns).to_proc updater = MassUpdater.new(@klass, key_col_proc: key_col_proc, col_proc: col_proc, replace_associations: self.replace_associations, bulk_import_id_generator: self.bulk_import_id_generator, save_associations_first: self.save_associations_first, bulk_import_id_column: self.bulk_import_id_column) saved, failures = save_record_list(record_list, updater) Karafka.monitor.instrument('deimos.batch_consumption.valid_records', { records: saved, consumer: self.class }) failures end |