Module: Datadog::Tracing::Contrib::Karafka::MessagesPatch

Defined in:
lib/datadog/tracing/contrib/karafka/patcher.rb

Overview

Patch to add tracing to Karafka::Messages::Messages

Instance Method Summary collapse

Instance Method Details

#configurationObject



13
14
15
# File 'lib/datadog/tracing/contrib/karafka/patcher.rb', line 13

def configuration
  Datadog.configuration.tracing[:karafka]
end

#each(&block) ⇒ Object

each is the most popular access point to Karafka messages, but not the only one Other access patterns do not have a straightforward tracing avenue (e.g. my_batch_operation messages.payloads)



26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
# File 'lib/datadog/tracing/contrib/karafka/patcher.rb', line 26

def each(&block)
  @messages_array.each do |message|
    trace_digest = if configuration[:distributed_tracing]
      headers = if message..respond_to?(:raw_headers)
        message..raw_headers
      else
        message..headers
      end
      Karafka.extract(headers)
    end

    Tracing.trace(Ext::SPAN_MESSAGE_CONSUME, continue_from: trace_digest) do |span, trace|
      if Datadog::DataStreams.enabled?
        begin
          headers = if message..respond_to?(:raw_headers)
            message..raw_headers
          else
            message..headers
          end

          Datadog::DataStreams.set_consume_checkpoint(
            type: "kafka",
            source: message.topic,
            auto_instrumentation: true
          ) { |key| headers[key] }
        rescue => e
          Datadog.logger.debug("Error setting DSM checkpoint: #{e.class}: #{e.message}")
        end
      end

      span.set_tag(Ext::TAG_OFFSET, message..offset)
      span.set_tag(Contrib::Ext::Messaging::TAG_DESTINATION, message.topic)
      span.set_tag(Contrib::Ext::Messaging::TAG_SYSTEM, Ext::TAG_SYSTEM)

      span.resource = message.topic

      yield message
    end
  end
end

#propagationObject



17
18
19
# File 'lib/datadog/tracing/contrib/karafka/patcher.rb', line 17

def propagation
  @propagation ||= Contrib::Karafka::Distributed::Propagation.new
end