影刀RPA店群自动化可观测性建设:日志、指标、链路追踪三位一体

在这里插入图片描述

系统跑起来之后,最怕的是什么?
不是出问题,而是出了问题不知道从哪里查。
我们早期排查一次故障的流程是这样的:先看哪个店铺报错,然后去服务器上grep日志,发现日志分散在十几个文件里,时间戳对不上。
再去看监控,发现只有CPU和内存,没有业务指标。
花了一个多小时才定位到是某个RPA流程里一个元素定位器失效了。

拼多多店群自动化报活动上架!


在这里插入图片描述

真正的问题不是没有日志,而是日志、指标、链路三者割裂,无法快速关联。
后来我们按照可观测性的三个支柱,重新构建了整个系统的可观测能力。
这篇文章就讲我们如何把日志、指标、链路追踪串联起来,让排障从小时级降到分钟级。

核心组件:结构化日志(JSON)、Prometheus指标、Jaeger链路追踪、Grafana统一大盘。


在这里插入图片描述

TEMU店群矩阵自动化运营核价报活动

一、可观测性的三个支柱

传统的监控只告诉你“系统出问题了”。
可观测性要回答三个问题:

  • 日志:具体发生了什么?(What happened?)
    • 指标:发生了多少次?趋势如何?(How many?)
    • 链路:在哪个环节出的问题?(Where exactly?)
      我们之前只做了指标(CPU/内存),日志是散乱的文本,链路完全没有。
      改造的第一步:统一日志格式,注入链路ID。

在这里插入图片描述

二、结构化日志:让日志可检索

以前的日志是这种格式:
在这里插入图片描述

  2024-01-15 10:23:45 INFO: Task started for shop pdd_001
    2024-01-15 10:23:47 ERROR: Element not found: #sync-btn
      ```
想查某个店铺的所有日志,只能grep。想按错误类型聚合,做不到。
我们改成了**JSON格式**,每条日志都是一个结构化对象。
```python
  # structured_logger.py
    import json
      import logging
        import sys
          from datetime import datetime
            from contextvars import ContextVar
trace_id_var = ContextVar('trace_id', default=None)
  shop_id_var = ContextVar('shop_id', default=None)
class JSONFormatter(logging.Formatter):
      def format(self, record):
                log_entry = {
                              "timestamp": datetime.utcnow().isoformat() + "Z",
                                            "level": record.levelname,
                                                          "logger": record.name,
                                                                        "message": record.getMessage(),
                                                                                      "trace_id": trace_id_var.get(),
                                                                                                    "shop_id": shop_id_var.get(),
                                                                                                                  "module": record.module,
                                                                                                                                "line": record.lineno
                                                                                                                                          }
                                                                                                                                                    if hasattr(record, 'exc_info') and record.exc_info:
                                                                                                                                                                  log_entry["exception"] = self.formatException(record.exc_info)
                                                                                                                                                                            return json.dumps(log_entry)
def setup_logging():
      handler = logging.StreamHandler(sys.stdout)
            handler.setFormatter(JSONFormatter())
                  root = logging.getLogger()
                        root.addHandler(handler)
                              root.setLevel(logging.INFO)
                                ```
每条日志自动带上 `trace_id` 和 `shop_id`。
查询时,在Kibana或Loki里直接按 `shop_id` 过滤,就能看到该店铺的所有相关日志。
按 `level: "ERROR"` 聚合,就能看到错误分布。
**我们当时在线上环境里踩过一次很严重的日志性能问题。**
JSON序列化每条日志有开销,在高并发下(每秒几百条日志)CPU占用达到15%。
解决方案:使用 `orjson` 替代标准 `json`,速度提升3倍。同时使用异步日志处理器,不阻塞主线程。
```python
  import orjson
    def orjson_serializer(record):
          return orjson.dumps(log_entry).decode()
            ```
---
## 三、指标:不仅仅是系统指标
我们之前只采集了系统指标(CPU、内存、磁盘)。
后来加上了**业务指标**和**RPA指标**。
```python
  # metrics_collector.py
    from prometheus_client import Counter, Histogram, Gauge, Info
      import time
class RPAMetrics:
      # 计数器
            task_total = Counter('rpa_task_total', 'Total tasks', ['platform', 'task_type', 'status'])
                  login_failures = Counter('rpa_login_failures', 'Login failures', ['platform'])
                        
                              # 直方图(耗时分布)
                                    task_duration = Histogram('rpa_task_duration_seconds', 'Task duration', ['task_type'], 
                                                                    buckets=[1, 2, 5, 10, 20, 30, 60, 120])
                                                                          
                                                                                # 仪表盘(当前值)
                                                                                      active_browsers = Gauge('rpa_active_browsers', 'Currently active browser instances')
                                                                                            queue_length = Gauge('rpa_queue_length', 'Task queue length', ['priority'])
                                                                                                  
                                                                                                        # 信息
                                                                                                              version_info = Info('rpa_version', 'System version')
                                                                                                                    
                                                                                                                          @classmethod
                                                                                                                                def record_task(cls, task_type, platform, status, duration_seconds):
                                                                                                                                          cls.task_total.labels(platform=platform, task_type=task_type, status=status).inc()
                                                                                                                                                    cls.task_duration.labels(task_type=task_type).observe(duration_seconds)
                                                                                                                                                      ```
这些指标暴露在 `/metrics` 端点,Prometheus每15秒抓取一次。
Grafana大盘上,我们可以实时看到:
- 每分钟完成任务数(按平台、任务类型)
-   - 任务P95耗时趋势
-   - 各平台登录失败率
-   - 浏览器实例池使用率
**一个有用的指标:店铺健康分**
我们把多个指标加权计算,得出每个店铺的健康分。
```python
  def calculate_shop_health(shop_id):
        login_success_rate = get_login_success_rate(shop_id, last_24h)
              task_success_rate = get_task_success_rate(shop_id, last_24h)
                    avg_task_duration = get_avg_task_duration(shop_id, last_24h)
                          
                                score = (
                                          login_success_rate * 0.4 +
                                                    task_success_rate * 0.4 +
                                                              max(0, 1 - avg_task_duration / 60) * 0.2
                                                                    ) * 100
                                                                          return score
                                                                            ```
健康分低于60的店铺自动标红,进入巡检队列。
---
## 四、链路追踪:看清调用链
RPA任务从入队到完成,会经过多个组件:
编排器 -> 消息队列 -> 执行节点 -> 浏览器实例池 -> 影刀RPA -> 数据清洗 -> 数据库
任何一个环节慢,都需要知道是哪里。
我们引入了OpenTelemetry,实现了分布式追踪。
```python
  # tracing_setup.py
    from opentelemetry import trace
      from opentelemetry.exporter.jaeger.thrift import JaegerExporter
        from opentelemetry.sdk.trace import TracerProvider
          from opentelemetry.sdk.trace.export import BatchSpanProcessor
            from opentelemetry.instrumentation.requests import RequestsInstrumentor
def init_tracing(service_name):
      provider = TracerProvider()
            jaeger_exporter = JaegerExporter(
                      agent_host_name="jaeger-agent",
                                agent_port=6831,
                                      )
                                            provider.add_span_processor(BatchSpanProcessor(jaeger_exporter))
                                                  trace.set_tracer_provider(provider)
                                                        tracer = trace.get_tracer(__name__)
                                                              RequestsInstrumentor().instrument()
                                                                    return tracer
                                                                      ```
在关键路径上手动添加Span:
```python
  tracer = init_tracing("rpa-executor")
def execute_task(task):
      with tracer.start_as_current_span("execute_task") as span:
                span.set_attribute("task_id", task["task_id"])
                          span.set_attribute("shop_id", task["shop_id"])
                                    
                                              with tracer.start_as_current_span("acquire_browser"):
                                                            browser = pool.acquire()
                                                                      
                                                                                with tracer.start_as_current_span("rpa_execute"):
                                                                                              result = call_rpa(browser, task)
                                                                                                        
                                                                                                                  with tracer.start_as_current_span("cleanup"):
                                                                                                                                pool.release(browser)
                                                                                                                                          
                                                                                                                                                    return result
                                                                                                                                                      ```
在Jaeger界面上,可以看到每个任务的完整调用链,以及每个环节的耗时。
有一次我们发现某个任务特别慢,打开链路发现瓶颈不在RPA,而在数据清洗时的数据库写入。原来是索引失效了。
没有链路追踪,这个问题可能要查半天。
---
## 五、日志、指标、链路的关联
三者不能各自孤立。
我们在日志中打印 `trace_id`,在指标中也可以带上 `trace_id` 的标签(但指标维度不能太高基数,所以只在大盘级别聚合)。
关键是在排障时,可以从指标异常点切入,找到对应的 `trace_id`,再查日志细节。
我们做了一个简单的**排障工作流**:
1. Grafana上看到任务失败率飙升
2.   2. 点击该时间段的图表,钻取到具体的失败任务列表
3.   3. 拿到失败任务的 `trace_id`
4.   4. 在Loki中搜索 `trace_id`,看到详细的错误日志
5.   5. 在Jaeger中查看该 `trace_id` 的调用链,定位到具体步骤
整个过程不到2分钟。
---
## 六、告警的工程化
有了指标,就可以配告警。
我们使用Prometheus的Alertmanager,配置了多级告警:
```yaml
  groups:
    - name: rpa_alerts
    -     rules:
    -     - alert: HighTaskFailureRate
    -       expr: |
    -         (sum(rate(rpa_task_total{status="failed"}[5m])) / 
    -          sum(rate(rpa_task_total[5m]))) > 0.1
    -       for: 5m
    -       annotations:
    -         summary: "任务失败率超过10%"
    -         
    -     - alert: QueueBacklog
    -       expr: rpa_queue_length > 500
    -       for: 2m
    -       annotations:
    -         summary: "队列积压超过500"
    -         
    -     - alert: LoginFailureSpike
    -       expr: increase(rpa_login_failures[10m]) > 10
    -       annotations:
    -         summary: "10分钟内登录失败超过10次"
    -   ```
告警路由到飞书,不同级别发送到不同群:
- P0(紧急):电话+飞书@all
-   - P1(严重):飞书@owner
-   - P2(警告):飞书群消息
我们还做了**告警收敛**,防止告警风暴。
同一个告警规则在30分钟内最多发送一次,并且合并相同店铺的告警。
```python
  # alert_aggregator.py
    class AlertAggregator:
          def __init__(self, redis_client):
                    self.redis = redis_client
                          
                                def should_send(self, alert_key, cooldown_minutes=30):
                                          last = self.redis.get(f"alert_last:{alert_key}")
                                                    if last and time.time() - float(last) < cooldown_minutes * 60:
                                                                  return False
                                                                            self.redis.setex(f"alert_last:{alert_key}", cooldown_minutes * 60, time.time())
                                                                                      return True
                                                                                        ```
---
## 七、可视化大盘设计
我们的Grafana大盘分为三层:
**第一层:全局概览**
- 总任务吞吐量(每分钟)
-   - 整体成功率(近24小时趋势)
-   - 当前队列长度
-   - 活跃浏览器实例数
-   - 各平台健康分
**第二层:平台详情**
- 拼多多/TEMU/TikTok Shop的独立指标
-   - 各任务类型的P95耗时
-   - 登录态分布(新鲜/即将过期/已过期)
**第三层:店铺下钻**
- 单个店铺的任务历史
-   - 店铺健康分变化曲线
-   - 最近失败记录
每个图表都支持点击下钻,直接跳转到对应日志或链路。
---
## 八、真实踩过的坑
**1. trace_id 在子进程中的传递**
我们使用 `multiprocessing` 启动子进程,`contextvars` 不会自动传递。
解决方法:在任务对象中显式传递 `trace_id`,子进程启动时重新设置。
```python
  def worker_main(task):
        trace_id = task["trace_id"]
              trace_id_var.set(trace_id)
                    # 继续执行
                      ```
**2. Prometheus 指标高基数问题**
一开始我们把 `shop_id` 作为指标标签,导致活跃序列数暴增(几千个店铺)。
Prometheus内存爆炸。后来去掉了 `shop_id` 标签,改用聚合指标。
店铺级别的指标通过其他途径(如日志)分析。
**3. Jaeger 存储占用太大**
全量采样导致 Jaeger 磁盘每天增长几十GB。
改用**概率采样**和**错误采样**:正常任务1%采样,错误任务100%采样。
```python
  from opentelemetry.sdk.trace.sampling import ParentBased, TraceIdRatioBased
    sampler = ParentBased(TraceIdRatioBased(0.01))
      ```
**4. 日志时区混乱**
容器内默认UTC,日志时间戳和北京时间差8小时,查问题很别扭。
统一在日志中使用UTC,展示时由前端转换。只要全局一致,不影响排查。
**5. 告警噪音太多**
初期规则太敏感,半夜经常被叫醒。后来加了 `for: 5m` 持续观察期,并引入静默时段(凌晨2-6点降级为消息)。
---
## 九、成果与收益
这套可观测性体系上线后,MTTR(平均修复时间)从47分钟降到了12分钟。
有一次凌晨3点,队列积压告警触发,我们打开大盘发现是某个TEMU店铺的同步任务卡住了,trace显示在登录环节超时。
远程重启了该店铺的浏览器实例,5分钟内恢复。
如果没有可观测性,可能要等到早上运营投诉才知道。
---
## 十、总结
可观测性不是可有可无的附加品,而是生产级系统的必需品。
建设的顺序建议:
6. 先做结构化日志 + trace_id 注入(最简单的改造)
7.   2. 加入关键业务指标(任务成功率、耗时)
8.   3. 用Prometheus + Grafana做可视化
9.   4. 接入链路追踪(只采样,慢慢调优)
10.  5. 配置分级告警和收敛
每一步都能带来可见的排障效率提升。
切忌一上来就想做成完美系统,从最痛的痛点开始。
希望这篇文章能帮你构建自己的可观测性体系。
---
作者:林焱 

更多推荐