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
66
|
# File 'lib/datadog/tracing/contrib/karafka/patcher.rb', line 26
def each(&block)
@messages_array.each do |message|
if configuration[:distributed_tracing]
= if message.metadata.respond_to?(:raw_headers)
message.metadata.
else
message.metadata.
end
trace_digest = Karafka.()
Datadog::Tracing.continue_trace!(trace_digest) if trace_digest
end
if Datadog::DataStreams.enabled?
begin
= if message.metadata.respond_to?(:raw_headers)
message.metadata.
else
message.metadata.
end
Datadog::DataStreams.set_consume_checkpoint(
type: 'kafka',
source: message.topic,
auto_instrumentation: true
) { |key| [key] }
rescue => e
Datadog.logger.debug("Error setting DSM checkpoint: #{e.class}: #{e}")
end
end
Tracing.trace(Ext::SPAN_MESSAGE_CONSUME) do |span|
span.set_tag(Ext::TAG_OFFSET, message.metadata.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
|