[Feature] Propagate Trace Headers into Root Span for OpenTelemetry Cross-Service Context (#10808)
Signed-off-by: zhanghaotong <zhanghaotong.zht@antgroup.com>
This commit is contained in:
@@ -82,6 +82,7 @@ from sglang.srt.server_args import (
|
|||||||
)
|
)
|
||||||
from sglang.srt.speculative.spec_info import SpeculativeAlgorithm
|
from sglang.srt.speculative.spec_info import SpeculativeAlgorithm
|
||||||
from sglang.srt.tracing.trace import (
|
from sglang.srt.tracing.trace import (
|
||||||
|
extract_trace_headers,
|
||||||
trace_get_proc_propagate_context,
|
trace_get_proc_propagate_context,
|
||||||
trace_req_finish,
|
trace_req_finish,
|
||||||
trace_req_start,
|
trace_req_start,
|
||||||
@@ -411,14 +412,18 @@ class TokenizerManager(TokenizerCommunicatorMixin):
|
|||||||
self.auto_create_handle_loop()
|
self.auto_create_handle_loop()
|
||||||
obj.normalize_batch_and_arguments()
|
obj.normalize_batch_and_arguments()
|
||||||
|
|
||||||
if request and "trace_context" in request.headers:
|
external_trace_header = None
|
||||||
trace_set_remote_propagate_context(request.headers["trace_context"])
|
if request:
|
||||||
|
if "trace_context" in request.headers:
|
||||||
|
trace_set_remote_propagate_context(request.headers["trace_context"])
|
||||||
|
else:
|
||||||
|
external_trace_header = extract_trace_headers(request.headers)
|
||||||
|
|
||||||
if self.server_args.tokenizer_worker_num > 1:
|
if self.server_args.tokenizer_worker_num > 1:
|
||||||
self._attach_multi_http_worker_info(obj)
|
self._attach_multi_http_worker_info(obj)
|
||||||
|
|
||||||
if self.enable_trace:
|
if self.enable_trace:
|
||||||
self._trace_request_start(obj, created_time)
|
self._trace_request_start(obj, created_time, external_trace_header)
|
||||||
|
|
||||||
if self.log_requests:
|
if self.log_requests:
|
||||||
max_length, skip_names, _ = self.log_request_metadata
|
max_length, skip_names, _ = self.log_request_metadata
|
||||||
@@ -2301,6 +2306,7 @@ class TokenizerManager(TokenizerCommunicatorMixin):
|
|||||||
self,
|
self,
|
||||||
obj: Union[GenerateReqInput, EmbeddingReqInput],
|
obj: Union[GenerateReqInput, EmbeddingReqInput],
|
||||||
created_time: Optional[float] = None,
|
created_time: Optional[float] = None,
|
||||||
|
external_trace_header: Optional[Dict] = None,
|
||||||
):
|
):
|
||||||
if obj.is_single:
|
if obj.is_single:
|
||||||
bootstrap_room = (
|
bootstrap_room = (
|
||||||
@@ -2311,6 +2317,7 @@ class TokenizerManager(TokenizerCommunicatorMixin):
|
|||||||
bootstrap_room,
|
bootstrap_room,
|
||||||
ts=int(created_time * 1e9),
|
ts=int(created_time * 1e9),
|
||||||
role=self.server_args.disaggregation_mode,
|
role=self.server_args.disaggregation_mode,
|
||||||
|
external_trace_header=external_trace_header,
|
||||||
)
|
)
|
||||||
trace_slice_start("", obj.rid, ts=int(created_time * 1e9), anonymous=True)
|
trace_slice_start("", obj.rid, ts=int(created_time * 1e9), anonymous=True)
|
||||||
else:
|
else:
|
||||||
@@ -2325,6 +2332,7 @@ class TokenizerManager(TokenizerCommunicatorMixin):
|
|||||||
bootstrap_room,
|
bootstrap_room,
|
||||||
ts=int(created_time * 1e9),
|
ts=int(created_time * 1e9),
|
||||||
role=self.server_args.disaggregation_mode,
|
role=self.server_args.disaggregation_mode,
|
||||||
|
external_trace_header=external_trace_header,
|
||||||
)
|
)
|
||||||
trace_slice_start(
|
trace_slice_start(
|
||||||
"", obj.rid[i], ts=int(created_time * 1e9), anonymous=True
|
"", obj.rid[i], ts=int(created_time * 1e9), anonymous=True
|
||||||
|
|||||||
@@ -30,10 +30,14 @@ from sglang.srt.utils import get_int_env_var
|
|||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from sglang.srt.managers.scheduler import Req
|
from sglang.srt.managers.scheduler import Req
|
||||||
|
from typing import Any, Dict, List, Mapping, Optional
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
opentelemetry_imported = False
|
opentelemetry_imported = False
|
||||||
tracing_enabled = False
|
tracing_enabled = False
|
||||||
|
_trace_context_propagator = None
|
||||||
|
|
||||||
|
TRACE_HEADERS = ["traceparent", "tracestate"]
|
||||||
|
|
||||||
try:
|
try:
|
||||||
from opentelemetry import context, propagate, trace
|
from opentelemetry import context, propagate, trace
|
||||||
@@ -49,6 +53,11 @@ try:
|
|||||||
from opentelemetry.sdk.resources import SERVICE_NAME, Resource
|
from opentelemetry.sdk.resources import SERVICE_NAME, Resource
|
||||||
from opentelemetry.sdk.trace import TracerProvider, id_generator
|
from opentelemetry.sdk.trace import TracerProvider, id_generator
|
||||||
from opentelemetry.sdk.trace.export import BatchSpanProcessor
|
from opentelemetry.sdk.trace.export import BatchSpanProcessor
|
||||||
|
from opentelemetry.trace.propagation.tracecontext import (
|
||||||
|
TraceContextTextMapPropagator,
|
||||||
|
)
|
||||||
|
|
||||||
|
_trace_context_propagator = TraceContextTextMapPropagator()
|
||||||
|
|
||||||
opentelemetry_imported = True
|
opentelemetry_imported = True
|
||||||
except ImportError:
|
except ImportError:
|
||||||
@@ -60,6 +69,14 @@ except ImportError:
|
|||||||
logger.info("opentelemetry package is not installed, tracing disabled")
|
logger.info("opentelemetry package is not installed, tracing disabled")
|
||||||
|
|
||||||
|
|
||||||
|
def is_tracing_enabled() -> bool:
|
||||||
|
return tracing_enabled
|
||||||
|
|
||||||
|
|
||||||
|
def extract_trace_headers(headers: Mapping[str, str]) -> Optional[Dict]:
|
||||||
|
return {h: headers[h] for h in TRACE_HEADERS if h in headers}
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class SglangTraceThreadInfo:
|
class SglangTraceThreadInfo:
|
||||||
host_id: str
|
host_id: str
|
||||||
@@ -418,6 +435,7 @@ def trace_req_start(
|
|||||||
bootstrap_room: Optional[int] = None,
|
bootstrap_room: Optional[int] = None,
|
||||||
ts: Optional[int] = None,
|
ts: Optional[int] = None,
|
||||||
role: Optional[str] = "null",
|
role: Optional[str] = "null",
|
||||||
|
external_trace_header: Optional[Dict[str, str]] = None,
|
||||||
):
|
):
|
||||||
if not tracing_enabled:
|
if not tracing_enabled:
|
||||||
return
|
return
|
||||||
@@ -444,10 +462,14 @@ def trace_req_start(
|
|||||||
tracer = threads_info[pid].tracer
|
tracer = threads_info[pid].tracer
|
||||||
if str(bootstrap_room) not in remote_trace_contexts:
|
if str(bootstrap_room) not in remote_trace_contexts:
|
||||||
attrs = {"bootstrap_room": str(hex(bootstrap_room))}
|
attrs = {"bootstrap_room": str(hex(bootstrap_room))}
|
||||||
|
external_trace_context = _trace_context_propagator.extract(
|
||||||
|
external_trace_header
|
||||||
|
)
|
||||||
bootstrap_room_span = tracer.start_span(
|
bootstrap_room_span = tracer.start_span(
|
||||||
name=f"Bootstrap Room {hex(bootstrap_room)}",
|
name=f"Bootstrap Room {hex(bootstrap_room)}",
|
||||||
start_time=ts,
|
start_time=ts,
|
||||||
attributes=attrs,
|
attributes=attrs,
|
||||||
|
context=external_trace_context,
|
||||||
)
|
)
|
||||||
reqs_context[rid].bootstrap_room_span = bootstrap_room_span
|
reqs_context[rid].bootstrap_room_span = bootstrap_room_span
|
||||||
bootstrap_room_span_context = trace.set_span_in_context(bootstrap_room_span)
|
bootstrap_room_span_context = trace.set_span_in_context(bootstrap_room_span)
|
||||||
|
|||||||
Reference in New Issue
Block a user