Class: Karafka::Pro::Processing::ConsumerGroups::Filters::Base

Inherits:
Object
  • Object
show all
Includes:
Core::Helpers::Time, Actions
Defined in:
lib/karafka/pro/processing/consumer_groups/filters/base.rb

Overview

Base for all the filters. All filters (including custom) need to use this API.

Due to the fact, that filters can limit data in such a way, that we need to pause or seek (throttling for example), the api is not just "remove some things from batch" but also provides ways to control the post-filtering operations that may be needed.

Constant Summary

Constants included from Actions

Actions::ALL

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from Actions

pause, #pause?, seek, #seek?, skip, #skip?

Constructor Details

#initializeBase

Initializes the filter as not yet applied



51
52
53
54
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 51

def initialize
  @applied = false
  @cursor = nil
end

Instance Attribute Details

#cursorKarafka::Messages::Message? (readonly)

Returns the message that we want to use as a cursor one to pause or seek or nil if not applicable.

Returns:

  • (Karafka::Messages::Message, nil)

    the message that we want to use as a cursor one to pause or seek or nil if not applicable.



44
45
46
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 44

def cursor
  @cursor
end

Instance Method Details

#actionSymbol

Returns filter post-execution action on consumer. One of Actions::ALL (Actions.skip, Actions.pause or Actions.seek).

Returns:

  • (Symbol)

    filter post-execution action on consumer. One of Actions::ALL (Actions.skip, Actions.pause or Actions.seek).



64
65
66
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 64

def action
  Actions.skip
end

#applied?Boolean

Returns did this filter change messages in any way.

Returns:

  • (Boolean)

    did this filter change messages in any way



69
70
71
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 69

def applied?
  @applied
end

#apply!(messages) ⇒ Object

Parameters:

  • messages (Array<Karafka::Messages::Message>)

    array with messages. Please keep in mind, this may already be partial due to execution of previous filters.

Raises:

  • (NotImplementedError)


58
59
60
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 58

def apply!(messages)
  raise NotImplementedError, "Implement in a subclass"
end

#mark_as_consumed?Boolean

Returns should we use the cursor value to mark as consumed. If any of the filters returns true, we return lowers applicable cursor value (if any).

Returns:

  • (Boolean)

    should we use the cursor value to mark as consumed. If any of the filters returns true, we return lowers applicable cursor value (if any)



82
83
84
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 82

def mark_as_consumed?
  false
end

#marking_cursorKarafka::Messages::Message?

Returns cursor message for marking or nil if no marking.

Returns:



94
95
96
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 94

def marking_cursor
  cursor
end

#marking_methodSymbol

Returns :mark_as_consumed or :mark_as_consumed!. Applicable only if marking is requested.

Returns:

  • (Symbol)

    :mark_as_consumed or :mark_as_consumed!. Applicable only if marking is requested



88
89
90
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 88

def marking_method
  :mark_as_consumed
end

#timeoutInteger?

Note:

Please do not return 0 when your filter is not pausing as it may interact with other filters that want to pause.

Returns default timeout for pausing (if applicable) or nil if not.

Returns:

  • (Integer, nil)

    default timeout for pausing (if applicable) or nil if not



76
77
78
# File 'lib/karafka/pro/processing/consumer_groups/filters/base.rb', line 76

def timeout
  nil
end