Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions lib/kafka/message_buffer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,40 @@ def each
end
end

def clear_largest_message
target = nil
@buffer.each do |topic, messages_for_topic|
messages_for_topic.each do |partition, messages_for_partition|
messages_for_partition.each do |message|
if target.nil? || target[:message].bytesize.to_i < message.bytesize.to_i
target = {
topic: topic,
partition: partition,
message: message,
}
end
end
end
end

clear_message(target) unless target.nil?
end

def clear_message(topic:, partition:, message:)
return if topic.nil? || partition.nil? || message.nil?
return unless @buffer.key?(topic) && @buffer[topic].key?(partition)
return unless @buffer[topic][partition].include?(message)

@size -= 1
@bytesize -= message.bytesize

@buffer[topic][partition].delete(message)
if @buffer[topic][partition].empty?
@buffer[topic].delete(partition)
@buffer.delete(topic) if @buffer[topic].empty?
end
end

# Clears buffered messages for the given topic and partition.
#
# @param topic [String] the name of the topic.
Expand Down
11 changes: 10 additions & 1 deletion lib/kafka/producer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -383,7 +383,16 @@ def deliver_messages_with_retries(notification)
end

assign_partitions!
operation.execute
begin
operation.execute
rescue Kafka::MessageSizeTooLarge => e
@logger.error "Received unrecoverable error #{e.class}: #{e.message}"
@logger.info "Dropping largest message from buffer. Current buffer size: #{@buffer.size}, bytesize: #{@buffer.bytesize}"
@buffer.clear_largest_message
@logger.info "Dropped largest message from buffer. Current buffer size: #{@buffer.size}, bytesize: #{@buffer.bytesize}"
sleep @retry_backoff
retry
end

if @required_acks.zero?
# No response is returned by the brokers, so we can't know which messages
Expand Down