SpanFromMessage reads the message
(ctx context.Context, tracer opentracing.Tracer, message *kafka.Message)
| 10 | |
| 11 | // SpanFromMessage reads the message |
| 12 | func 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 | |
| 31 | func getCarrier(msg *kafka.Message) opentracing.TextMapCarrier { |
| 32 | mapCarrier := make(opentracing.TextMapCarrier) |