Class: Deimos::Utils::DbProducer

Inherits:
Object
  • Object
show all
Includes:
Phobos::Producer
Defined in:
sig/defs.rbs

Overview

Class which continually polls the kafka_messages table in the database and sends Kafka messages.

Constant Summary collapse

BATCH_SIZE =

Returns:

  • (Integer)
DELETE_BATCH_SIZE =

Returns:

  • (Integer)
MAX_DELETE_ATTEMPTS =

Returns:

  • (Integer)

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initializeDbProducer

@param logger

Parameters:

  • logger (Logger)


816
# File 'sig/defs.rbs', line 816

def initialize: (?Logger logger) -> void

Instance Attribute Details

#current_topicObject

Returns the value of attribute current_topic.

Returns:

  • (Object)


868
869
870
# File 'sig/defs.rbs', line 868

def current_topic
  @current_topic
end

#idObject

Returns the value of attribute id.

Returns:

  • (Object)


865
866
867
# File 'sig/defs.rbs', line 865

def id
  @id
end

Instance Method Details

#compact_messages::Array[Deimos::KafkaMessage]

@param batch

Parameters:

Returns:



862
# File 'sig/defs.rbs', line 862

def compact_messages: (::Array[Deimos::KafkaMessage] batch) -> ::Array[Deimos::KafkaMessage]

#configFigTree

Returns:



818
# File 'sig/defs.rbs', line 818

def config: () -> FigTree

#delete_messagesvoid

This method returns an undefined value.

@param messages

Parameters:



840
# File 'sig/defs.rbs', line 840

def delete_messages: (::Array[Deimos::KafkaMessage] messages) -> void

#log_messagesvoid

This method returns an undefined value.

@param messages

Parameters:



845
# File 'sig/defs.rbs', line 845

def log_messages: (::Array[Deimos::KafkaMessage] messages) -> void

#process_next_messagesvoid

This method returns an undefined value.

Complete one loop of processing all messages in the DB.



827
# File 'sig/defs.rbs', line 827

def process_next_messages: () -> void

#process_topicString?

@param topic

@return — the topic that was locked, or nil if none were.

Parameters:

  • topic (String)

Returns:

  • (String, nil)


834
# File 'sig/defs.rbs', line 834

def process_topic: (String topic) -> String?

#process_topic_batchvoid

This method returns an undefined value.

Process a single batch in a topic.



837
# File 'sig/defs.rbs', line 837

def process_topic_batch: () -> void

#produce_messagesvoid

This method returns an undefined value.

Produce messages in batches, reducing the size 1/10 if the batch is too large. Does not retry batches of messages that have already been sent.

@param batch

Parameters:

  • batch (::Array[::Hash[untyped, untyped]])


859
# File 'sig/defs.rbs', line 859

def produce_messages: (::Array[::Hash[untyped, untyped]] batch) -> void

#retrieve_messages::Array[Deimos::KafkaMessage]

Returns:



842
# File 'sig/defs.rbs', line 842

def retrieve_messages: () -> ::Array[Deimos::KafkaMessage]

#retrieve_topics::Array[String]

Returns:

  • (::Array[String])


829
# File 'sig/defs.rbs', line 829

def retrieve_topics: () -> ::Array[String]

#send_pending_metricsvoid

This method returns an undefined value.

Send metrics related to pending messages.



848
# File 'sig/defs.rbs', line 848

def send_pending_metrics: () -> void

#shutdown_producervoid

This method returns an undefined value.

Shut down the sync producer if we have to. Phobos will automatically create a new one. We should call this if the producer can be in a bad state and e.g. we need to clear the buffer.



853
# File 'sig/defs.rbs', line 853

def shutdown_producer: () -> void

#startvoid

This method returns an undefined value.

Start the poll.



821
# File 'sig/defs.rbs', line 821

def start: () -> void

#stopvoid

This method returns an undefined value.

Stop the poll.



824
# File 'sig/defs.rbs', line 824

def stop: () -> void