ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

FastStream 集成 OpenTelemetry:为 Kafka / RabbitMQ / NATS / Redis / MQTT 服务构建统一可观测性

FastStream 集成 OpenTelemetry:为 Kafka / RabbitMQ / NATS / Redis / MQTT 服务构建统一可观测性 FastStream 集成 OpenTelemetry为 Kafka / RabbitMQ / NATS / Redis / MQTT 服务构建统一可观测性【免费下载链接】faststreamAsynchronous Python framework for event-driven services. A thin client for Kafka, RabbitMQ, NATS, Redis and MQTT with full access to native broker features, plus AsyncAPI docs, in-memory tests and observability out of the box.项目地址: https://gitcode.com/GitHub_Trending/fa/faststreamFastStream 是面向事件驱动服务的异步 Python 框架其内置的 OpenTelemetry 支持可以让开发者在 Kafka、RabbitMQ、NATS、Redis 与 MQTT 等异构消息系统上获得统一且开箱即用的链路追踪Traces与指标Metrics能力。本文以仓库中的 OpenTelemetry 官方指南docs/docs/en/getting-started/observability/opentelemetry/index.md为骨架结合 faststream/opentelemetry 的中间件源码与 tests/opentelemetry 测试用例讲解从零接入、Span 模型、指标采集到 Baggage 上下文传递的完整实战方案。读完本文你将能在一个 FastStream 服务上快速启用追踪与指标并理解其底层传播与关联机制。什么是 OpenTelemetryFastStream 为何内置支持它OpenTelemetry是一个开源的可观测性框架旨在为 traces链路追踪、metrics指标和 logs日志的采集与导出提供统一标准让可观测性成为软件开发的内置能力从而简化不同服务之间遥测数据的集成与标准化。FastStream 之所以把 OpenTelemetry 作为一等公民内置是因为消息驱动的微服务天然跨服务、跨 Broker追踪上下文需要在生产者 → Broker → 消费者整条链路上无损传递。从 pyproject.toml 可以看到官方通过faststream[otel]可选依赖引入opentelemetry-sdk1.24.0,2.0.0并在此基础上实现了通用遥测中间件 faststream/opentelemetry/middleware.py各 Broker 专属的TelemetryMiddleware与属性提供器如 faststream/kafka/opentelemetry方便在业务处理器中直接取用的CurrentSpan/CurrentBaggage注解与Baggage工具类见 faststream/opentelemetry/annotations.py、faststream/opentelemetry/baggage.py。因此你无需为每个消息系统编写单独的埋点代码只需一行中间件即可让所有订阅与发布动作自动产生符合语义约定Semantic Conventions的遥测数据。快速开始三步为 FastStream 服务接入 OpenTelemetry官方文档给出的接入流程非常精简只需三个步骤。第一步安装依赖pip install faststream[otel] opentelemetry-exporter-otlp其中faststream[otel]会安装opentelemetry-sdk及其依赖opentelemetry-exporter-otlp负责将遥测数据通过 OTLP 协议导出到后端如 Collector、Jaeger、Grafana Tempo 等。第二步配置 TracerProvider 与 gRPC 导出器from opentelemetry import trace from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import TracerProvider from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter from opentelemetry.sdk.trace.export import BatchSpanProcessor resource Resource.create(attributes{service.name: faststream}) tracer_provider TracerProvider(resourceresource) trace.set_tracer_provider(tracer_provider) exporter OTLPSpanExporter(endpointhttp://localhost:4317) processor BatchSpanProcessor(exporter) tracer_provider.add_span_processor(processor)要点说明service.name属性用于标识当前服务在后端中会作为服务名的维度展示OTLPSpanExporter默认通过 gRPC端口4317导出如需 HTTP/protobuf 方式可换成对应的OTLPSpanExporter变体BatchSpanProcessor会批量发送 Span降低网络开销测试环境也可改用SimpleSpanProcessor即时导出仓库测试即采用该方式见 tests/opentelemetry/basic.py。第三步为 Broker 挂载遥测中间件不同消息系统的导入路径与中间件名称不同但用法完全一致——在创建 Broker 时传入middlewares元组即可代码均收录于 docs/docs_src/getting_started/opentelemetryAIOKafkakafka_telemetry.pyfrom faststream import FastStream from faststream.kafka import KafkaBroker from faststream.kafka.opentelemetry import KafkaTelemetryMiddleware broker KafkaBroker( middlewares( KafkaTelemetryMiddleware(tracer_providertracer_provider), ), ) app FastStream(broker)Confluent Kafkaconfluent_telemetry.pyfrom faststream import FastStream from faststream.confluent import KafkaBroker from faststream.confluent.opentelemetry import KafkaTelemetryMiddleware broker KafkaBroker( middlewares( KafkaTelemetryMiddleware(tracer_providertracer_provider), ), ) app FastStream(broker)RabbitMQrabbit_telemetry.pyfrom faststream.rabbit import RabbitBroker from faststream.rabbit.opentelemetry import RabbitTelemetryMiddleware broker RabbitBroker( middlewares( RabbitTelemetryMiddleware(tracer_providertracer_provider), ), ) app FastStream(broker)NATSnats_telemetry.pyfrom faststream.nats import NatsBroker from faststream.nats.opentelemetry import NatsTelemetryMiddleware broker NatsBroker( middlewares( NatsTelemetryMiddleware(tracer_providertracer_provider), ), ) app FastStream(broker)Redisredis_telemetry.pyfrom faststream.redis import RedisBroker from faststream.redis.opentelemetry import RedisTelemetryMiddleware broker RedisBroker( middlewares( RedisTelemetryMiddleware(tracer_providertracer_provider), ), ) app FastStream(broker)MQTTmqtt_telemetry.pyfrom faststream.mqtt import MQTTBroker from faststream.mqtt.opentelemetry import MQTTTelemetryMiddleware broker MQTTBroker( middlewares( MQTTTelemetryMiddleware(tracer_providertracer_provider), ), ) app FastStream(broker)MQTT 兼容性警告OpenTelemetry 中间件仅支持MQTT 5不兼容 MQTT 3.1.1——因为 3.1.1 协议不支持用于 trace context 传播的 user properties用户属性。中间件参数详解各 Broker 的TelemetryMiddleware构造函数在底层统一收敛到 faststream/opentelemetry/middleware.py 的TelemetryMiddleware.__init__常用参数如下参数默认值作用tracer_providerNone使用全局 Provider指定TracerProvider未传时通过trace.get_tracer(...)获取全局 Tracermeter_providerNone指定MeterProvider用于创建指标meterNone直接传入自定义Meter优先级高于meter_providerinclude_messages_counters基础类为FalseKafka 等为True是否额外记录消息条数计数器见下文指标一节以 faststream/kafka/opentelemetry/middleware.py 为例KafkaTelemetryMiddleware默认include_messages_countersTrue并通过telemetry_attributes_provider_factory在单条消息与批量消息之间自动选择对应的属性提供器。深入原理FastStream 的 Span 模型与上下文传播三种 Spancreate / publish / process从 faststream/opentelemetry/consts.py 可以看到MessageAction定义了四个动作create、publish、process、receive。结合中间件源码middleware.py可以还原出完整链路发布侧若无当前活动 Span先以SpanKind.PRODUCER创建一个{destination} create的短 Span 记录消息被创建这一事件紧接着以SpanKind.PRODUCER创建{destination} publishSpan 包裹真实的发布调用消费侧从消息头中提取 W3C trace contextTraceContextTextMapPropagator若没有上游上下文则以SpanKind.CONSUMER创建{destination} create随后以SpanKind.CONSUMER创建{destination} processSpan 包裹处理器执行处理结束after_processed时关闭 Span。Span 命名规则为f{destination} {action}见 middleware.pydestination 即目标主题/队列/频道名。tests/opentelemetry/basic.py中的assert_span验证了 Span 名称、kind、父子关系与关键属性basic.py而test_subscriber_create_publish_process_span断言一次消息流转会产生create → publish → process三个 Spanbasic.py。Span 属性Attributes各 Broker 通过实现TelemetrySettingsProvider协议见 faststream/opentelemetry/provider.py提供消费/发布两侧的属性。以 Kafka 为例faststream/kafka/opentelemetry/provider.pySpan 上会携带通用属性messaging.system如kafka、messaging.destination.name发布目标、messaging.message.id消息 IDUUID、messaging.message.conversation_id关联 IDUUID、messaging.message.payload_size_bytes负载大小Kafka 专属属性messaging.kafka.destination.partition分区、messaging.kafka.message.offset偏移量、messaging.kafka.message.key消息键批量场景messaging.batch.message.count批量条数异常场景error.type异常类名见 consts.py。这些属性都遵循 OpenTelemetry 的 Semantic Conventions方便在 Grafana Tempo、Jaeger 等后端直接按消息系统、主题、分区做筛选。跨服务上下文传播Trace Context BaggageFastStream 使用标准 W3C 传播器middleware.pyTraceContextTextMapPropagator负责traceparent/tracestate的注入与提取保证分布式链路跨服务连接W3CBaggagePropagator负责 baggage 键值对的注入与提取用于跨服务携带业务上下文如用户 ID、租户 ID。发布时publish_scope中间件把当前活动 Span 写入上下文并注入到消息头中消费时consume_scope从消息头提取上下文作为processSpan 的父上下文同时把 Span 与 Baggage 写入 FastStream 的 ContextRepo从而让处理器内任意位置都能访问middleware.py。在处理器中访问当前 Span 与 Baggage通过 faststream/opentelemetry/annotations.py 提供的注解类型可以直接把当前 Span 和 Baggage 注入处理器参数from faststream.opentelemetry import CurrentBaggage, CurrentSpan broker.subscriber(orders) async def handler(msg: str, span: CurrentSpan, baggage: CurrentBaggage) - None: # span 即当前 process Span可直接添加自定义属性 span.set_attribute(handler.name, orders.handler) # baggage 支持 get / set / remove / clear baggage.set(tenant_id, t-001) ...Baggage类faststream/opentelemetry/baggage.py提供get/set/remove/clear/get_all/to_headers等方法Baggage.from_msg还支持从批量消息中合并多个子消息的 baggage。测试test_get_baggage、test_modify_baggage、test_clear_baggage分别验证了读取、修改、清除与跨处理器传播的行为basic.py。批量消息与 Span Link对于批量消费FastStream 不会强行生成父子关系而是为processSpan 关联links从每条子消息的 header 提取各自的上游 Span 上下文若同一correlation_id出现多次还会带上messaging.batch.message.count属性标注重复计数middleware.py。同时通过 baggage 中的with_batch标记识别批量消息consts.py。指标Metrics消息处理的时长与计数除了链路追踪中间件还内置了基于 OpenTelemetry Metrics 的指标采集见 middleware.py 的_MetricsContainer指标名类型单位说明messaging.publish.durationHistograms发布操作耗时messaging.process.durationHistograms消息处理耗时messaging.publish.messagesCountermessage发布消息条数需开启计数器messaging.process.messagesCountermessage处理消息条数需开启计数器指标维度attributes包括messaging.system、messaging.destination.name发布侧或messaging.destination_publish.name消费侧见 consts.py当处理/发布抛出异常时还会附加error.type维度标记错误类型middleware.py。在测试中test_metrics通过InMemoryMetricReader断言 4 个指标开启计数器时或 2 个指标关闭时均被正确记录test_error_metrics则验证了error.type维度basic.py。由于各 Broker 实现均基于同一_MetricsContainer因此所有消息系统暴露的指标口径完全一致可直接用 Prometheus 抓取后接入 Grafana 面板。参考示例一套完整的 FastStream 可观测性基础设施官方文档提供了一个预配置的参考项目可作为搭建服务与基础设施的范本示例为仓库文档中引用的faststream-monitoring示例工程。该示例包含三个 FastStream 服务演示跨服务链路通过gRPC将 traces 导出到Grafana Tempo使用Grafana可视化分布式追踪采集指标并通过Prometheus导出一套用于指标的Grafana dashboard自定义 Span 的使用示例一套预配置的docker-compose包含上述全部基础设施。其分布式追踪在 Grafana Tempo 中的可视化效果如下图所示——可以看到orders、trades、notifications三个服务之间的publish/processSpan 首尾相接、耗时一目了然结语FastStream 的 OpenTelemetry 集成把可观测性变成了框架的内置能力一条中间件即可覆盖 6 大消息系统自动生成语义规范的create / publish / processSpan、跨 Broker 的上下文传播、开箱即用的消息指标并支持在处理器内直接访问CurrentSpan与CurrentBaggage。如需进一步深入可继续阅读中间件核心实现faststream/opentelemetry/middleware.py各 Broker 属性提供器以 faststream/kafka/opentelemetry/provider.py 为参考同构目录见 faststream 下各 Broker 的opentelemetry子包官方代码示例docs/docs_src/getting_started/opentelemetry行为验证测试tests/opentelemetry每个 Broker 一个测试文件共享 basic.py 中的公共测试基类【免费下载链接】faststreamAsynchronous Python framework for event-driven services. A thin client for Kafka, RabbitMQ, NATS, Redis and MQTT with full access to native broker features, plus AsyncAPI docs, in-memory tests and observability out of the box.项目地址: https://gitcode.com/GitHub_Trending/fa/faststream创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表