Class: Karafka::Routing::Topic

Inherits:
Object
  • Object
show all
Defined in:
lib/karafka/routing/topic.rb

Overview

Topic stores all the details on how we should interact with Kafka given topic. It belongs to a consumer group as from 0.6 all the topics can work in the same consumer group It is a part of Karafka’s DSL.

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(name, consumer_group) ⇒ Topic

Returns a new instance of Topic.

Parameters:



35
36
37
38
39
40
41
42
43
44
# File 'lib/karafka/routing/topic.rb', line 35

def initialize(name, consumer_group)
  @name = name.to_s
  @consumer_group = consumer_group
  @attributes = {}
  @active = true
  # @note We use identifier related to the consumer group that owns a topic, because from
  #   Karafka 0.6 we can handle multiple Kafka instances with the same process and we can
  #   have same topic name across multiple consumer groups
  @id = "#{consumer_group.id}_#{@name}"
end

Instance Attribute Details

#consumerClass

Returns consumer class that we should use.

Returns:

  • (Class)

    consumer class that we should use



66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
# File 'lib/karafka/routing/topic.rb', line 66

def consumer
  if consumer_persistence
    # When persistence of consumers is on, no need to reload them
    @consumer
  else
    # In order to support code reload without having to change the topic api, we re-fetch the
    # class of a consumer based on its class name. This will support all the cases where the
    # consumer class is defined with a name. It won't support code reload for anonymous
    # consumer classes, but this is an edge case
    begin
      ::Object.const_get(@consumer.to_s)
    rescue NameError
      # It will only fail if the in case of anonymous classes
      @consumer
    end
  end
end

#consumer_groupObject (readonly)

Returns the value of attribute consumer_group.



9
10
11
# File 'lib/karafka/routing/topic.rb', line 9

def consumer_group
  @consumer_group
end

#idObject (readonly)

Returns the value of attribute id.



9
10
11
# File 'lib/karafka/routing/topic.rb', line 9

def id
  @id
end

#nameObject (readonly)

Returns the value of attribute name.



9
10
11
# File 'lib/karafka/routing/topic.rb', line 9

def name
  @name
end

#subscription_groupObject

Full subscription group reference can be built only when we have knowledge about the whole routing tree, this is why it is going to be set later on



16
17
18
# File 'lib/karafka/routing/topic.rb', line 16

def subscription_group
  @subscription_group
end

#subscription_group_nameObject

Returns the value of attribute subscription_group_name.



12
13
14
# File 'lib/karafka/routing/topic.rb', line 12

def subscription_group_name
  @subscription_group_name
end

Instance Method Details

#active(active) ⇒ Object

Allows to disable topic by invoking this method and setting it to ‘false`.

Parameters:

  • active (Boolean)

    should this topic be consumed or not



86
87
88
# File 'lib/karafka/routing/topic.rb', line 86

def active(active)
  @active = active
end

#active?Boolean

Returns should this topic be in use.

Returns:

  • (Boolean)

    should this topic be in use



100
101
102
103
104
105
# File 'lib/karafka/routing/topic.rb', line 100

def active?
  # Never active if disabled via routing
  return false unless @active

  Karafka::App.config.internal.routing.activity_manager.active?(:topics, name)
end

#consumer_classClass

Note:

This is just an alias to the ‘#consumer` method. We however want to use it internally instead of referencing the `#consumer`. We use this to indicate that this method returns class and not an instance. In the routing we want to keep the `#consumer Consumer` routing syntax, but for references outside, we should use this one.

Returns consumer class that we should use.

Returns:

  • (Class)

    consumer class that we should use



95
96
97
# File 'lib/karafka/routing/topic.rb', line 95

def consumer_class
  consumer
end

#subscription_nameString

Returns name of subscription that will go to librdkafka.

Returns:

  • (String)

    name of subscription that will go to librdkafka



61
62
63
# File 'lib/karafka/routing/topic.rb', line 61

def subscription_name
  name
end

#to_hHash

Note:

This is being used when we validate the consumer_group and its topics

Returns hash with all the topic attributes.

Returns:

  • (Hash)

    hash with all the topic attributes



109
110
111
112
113
114
115
116
117
118
119
120
121
122
# File 'lib/karafka/routing/topic.rb', line 109

def to_h
  map = INHERITABLE_ATTRIBUTES.map do |attribute|
    [attribute, public_send(attribute)]
  end

  Hash[map].merge!(
    id: id,
    name: name,
    active: active?,
    consumer: consumer,
    consumer_group_id: consumer_group.id,
    subscription_group_name: subscription_group_name
  ).freeze
end