Class: Karafka::Routing::Topic

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

Overview

Note:

#group is the polymorphic reference to the owning group. Today this is always a ConsumerGroup, but the accessor is named generically in preparation for additional group types (e.g. KIP-932 share groups). #consumer_group is kept as an alias for backwards compatibility.

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

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(name, group) ⇒ Topic

Returns a new instance of Topic.

Parameters:

  • name (String, Symbol)

    name of a topic on which we want to listen

  • group (Karafka::Routing::ConsumerGroup)

    owning group of this topic. Polymorphic placeholder for future group types (e.g. share groups); today always a ConsumerGroup.



42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
# File 'lib/karafka/routing/topic.rb', line 42

def initialize(name, group)
  @name = name.to_s
  @group = group
  @attributes = {}
  @active = true
  # We use identifier related to the 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 groups
  @id = "#{group.id}_#{@name}"
  @consumer = nil
  @active_assigned = false
  @subscription_group_details = nil

  INHERITABLE_ATTRIBUTES.each do |attribute|
    instance_variable_set("@#{attribute}", nil)
  end

  # Explicit nil initialization for Ruby's object shapes optimization. The per-topic pause
  # config is built lazily on first read, defaulting to the global `config.pause.*` settings.
  @pause = nil
end

Instance Attribute Details

#consumerClass

Returns consumer class that we should use.

Returns:

  • (Class)

    consumer class that we should use



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

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

#groupObject (readonly) Also known as: consumer_group

Returns the value of attribute group.



14
15
16
# File 'lib/karafka/routing/topic.rb', line 14

def group
  @group
end

#idObject (readonly)

Returns the value of attribute id.



14
15
16
# File 'lib/karafka/routing/topic.rb', line 14

def id
  @id
end

#nameObject (readonly)

Returns the value of attribute name.



14
15
16
# File 'lib/karafka/routing/topic.rb', line 14

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



26
27
28
# File 'lib/karafka/routing/topic.rb', line 26

def subscription_group
  @subscription_group
end

#subscription_group_detailsObject

Returns the value of attribute subscription_group_details.



22
23
24
# File 'lib/karafka/routing/topic.rb', line 22

def subscription_group_details
  @subscription_group_details
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



128
129
130
131
132
133
134
135
136
137
138
139
# File 'lib/karafka/routing/topic.rb', line 128

def active(active)
  # Do not allow for active overrides. Basically if this is set on the topic level, defaults
  # will not overwrite it and this is desired. Otherwise because of the fact that this is
  # not a full feature config but just a flag, default value would always overwrite the
  # per-topic config since defaults application happens after the topic config block
  unless @active_assigned
    @active = active
    @active_assigned = true
  end

  @active
end

#active?Boolean

Returns should this topic be in use.

Returns:

  • (Boolean)

    should this topic be in use



151
152
153
154
155
156
# File 'lib/karafka/routing/topic.rb', line 151

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



146
147
148
# File 'lib/karafka/routing/topic.rb', line 146

def consumer_class
  consumer
end

#kafka=(settings = {}) ⇒ Object

Note:

It is set to false by default to preserve backwards compatibility

Often users want to have the same basic cluster setup with small setting alterations This method allows us to do so by setting inherit to true. Whe inherit is enabled, settings will be merged with defaults.

Parameters:

  • settings (Hash) (defaults to: {})

    kafka scope settings. If :inherit key is provided, it will instruct the assignment to merge with root level defaults



96
97
98
99
100
# File 'lib/karafka/routing/topic.rb', line 96

def kafka=(settings = {})
  inherit = settings.delete(:inherit)

  @kafka = inherit ? Karafka::App.config.kafka.merge(settings) : settings
end

#pauseKarafka::Routing::Features::Pausing::Config

Returns per-topic pause configuration, reflecting the root config.pause.* settings.

Returns:



79
80
81
82
83
84
85
86
# File 'lib/karafka/routing/topic.rb', line 79

def pause
  @pause ||= Features::Pausing::Config.new(
    active: false,
    timeout: Karafka::App.config.pause.timeout,
    max_timeout: Karafka::App.config.pause.max_timeout,
    with_exponential_backoff: Karafka::App.config.pause.with_exponential_backoff
  )
end

#subscription_nameString

Returns name of subscription that will go to librdkafka.

Returns:

  • (String)

    name of subscription that will go to librdkafka



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

def subscription_name
  name
end

#to_hHash

Note:

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

Returns hash with all the topic attributes.

Returns:

  • (Hash)

    hash with all the topic attributes



160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
# File 'lib/karafka/routing/topic.rb', line 160

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

  map.to_h.merge!(
    id: id,
    name: name,
    active: active?,
    consumer: consumer,
    pause: pause.to_h,
    group_id: group.id,
    # Kept as a reference alongside `group_id` for backwards compatibility. Will be removed
    # in Karafka 3.0.
    consumer_group_id: group.id,
    subscription_group_details: subscription_group_details
  ).freeze
end