MCPcopy Create free account
hub / github.com/DoNewsCode/core / SpanFromMessage

Function SpanFromMessage

otkafka/helper.go:12–29  ·  view source on GitHub ↗

SpanFromMessage reads the message

(ctx context.Context, tracer opentracing.Tracer, message *kafka.Message)

Source from the content-addressed store, hash-verified

10
11// SpanFromMessage reads the message
12func SpanFromMessage(ctx context.Context, tracer opentracing.Tracer, message *kafka.Message) (opentracing.Span, context.Context, error) {
13 carrier := getCarrier(message)
14 spanContext, err := tracer.Extract(opentracing.TextMap, carrier)
15 if err != nil && err != opentracing.ErrSpanContextNotFound {
16 return nil, nil, err
17 }
18 span := tracer.StartSpan("kafka reader", ext.RPCServerOption(spanContext))
19 if err != nil && err != opentracing.ErrSpanContextNotFound {
20 return nil, nil, err
21 }
22 ext.SpanKind.Set(span, ext.SpanKindConsumerEnum)
23 ext.PeerService.Set(span, "kafka")
24 span.SetTag("topic", message.Topic)
25 span.SetTag("partition", message.Partition)
26 span.SetTag("offset", message.Offset)
27
28 return span, opentracing.ContextWithSpan(ctx, span), nil
29}
30
31func getCarrier(msg *kafka.Message) opentracing.TextMapCarrier {
32 mapCarrier := make(opentracing.TextMapCarrier)

Callers 2

TestHelper_no_parentFunction · 0.85
TestWriterFunction · 0.85

Calls 2

getCarrierFunction · 0.85
SetMethod · 0.45

Tested by 2

TestHelper_no_parentFunction · 0.68
TestWriterFunction · 0.68