返利软件后端日志体系构建:结构化日志 + ELK 链路追踪用于问题定位

大家好,我是 微赚淘客系统3.0 的研发者省赚客!

在高并发、多服务的返利系统中,用户一次“查券+下单+返利”操作可能涉及 5 个以上微服务。若出现佣金未到账或查询超时,传统文本日志难以快速定位根因。微赚淘客系统3.0 基于 JSON 结构化日志 + TraceID 全链路透传 + ELK(Elasticsearch + Logstash + Kibana) 构建可观测性体系,实现秒级问题定位。

一、日志格式标准化

所有服务强制输出 JSON 格式日志,包含关键字段:

{
  "timestamp": "2026-01-30T14:23:01.123Z",
  "level": "INFO",
  "service": "rebate-service",
  "traceId": "a1b2c3d4-e5f6-7890-g1h2-i3j4k5l6m7n8",
  "spanId": "s9o8p7q6",
  "userId": 10086,
  "tradeId": "T20260130142301",
  "event": "COMMISSION_CALCULATED",
  "message": "佣金计算成功",
  "amount": 12.50,
  "durationMs": 45
}

使用 Logback 配置 JSON 输出:

<!-- logback-spring.xml -->
<configuration>
  <appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
    <encoder class="net.logstash.logback.encoder.LoggingEventCompositeJsonEncoder">
      <providers>
        <timestamp/>
        <logLevel/>
        <loggerName>
          <fieldName>service</fieldName>
        </loggerName>
        <mdc/> <!-- 注入 traceId、userId 等 -->
        <message/>
        <arguments/>
        <stackTrace/>
      </providers>
    </encoder>
  </appender>

  <root level="INFO">
    <appender-ref ref="STDOUT"/>
  </root>
</configuration>

二、TraceID 全链路透传

基于 Spring Cloud Sleuth 自动注入 traceIdspanId,并通过 MDC 传递至日志:

package juwatech.cn.rebate.config;

import org.slf4j.MDC;
import org.springframework.cloud.sleuth.Tracer;
import org.springframework.stereotype.Component;

import javax.servlet.Filter;
import javax.servlet.FilterChain;
import javax.servlet.ServletException;
import javax.servlet.ServletRequest;
import javax.servlet.ServletResponse;
import java.io.IOException;

@Component
public class TraceIdMdcFilter implements Filter {

    private final Tracer tracer;

    public TraceIdMdcFilter(Tracer tracer) {
        this.tracer = tracer;
    }

    @Override
    public void doFilter(ServletRequest request, ServletResponse response,
                         FilterChain chain) throws IOException, ServletException {
        String traceId = tracer.currentSpan() != null ? 
            tracer.currentSpan().context().traceIdString() : "N/A";
        MDC.put("traceId", traceId);
        try {
            chain.doFilter(request, response);
        } finally {
            MDC.clear();
        }
    }
}

在业务代码中手动记录关键事件:

package juwatech.cn.rebate.service;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;

@Service
public class CommissionService {

    private static final Logger log = LoggerFactory.getLogger(CommissionService.class);

    public void processRebate(Long userId, String tradeId) {
        long start = System.currentTimeMillis();
        try {
            // 业务逻辑
            BigDecimal amount = calculateCommission(tradeId);
            
            log.info("{}", Map.of(
                "event", "COMMISSION_CALCULATED",
                "userId", userId,
                "tradeId", tradeId,
                "amount", amount,
                "durationMs", System.currentTimeMillis() - start
            ));
            
            accountService.credit(userId, amount);
            
            log.info("{}", Map.of(
                "event", "REBATE_SUCCESS",
                "userId", userId,
                "tradeId", tradeId
            ));
        } catch (Exception e) {
            log.error("{}", Map.of(
                "event", "REBATE_FAILED",
                "userId", userId,
                "tradeId", tradeId,
                "error", e.getMessage(),
                "stack", ExceptionUtils.getStackTrace(e)
            ));
            throw e;
        }
    }
}

三、Logstash 配置解析与索引

Logstash 从 Kafka 消费日志并写入 Elasticsearch:

# logstash.conf
input {
  kafka {
    bootstrap_servers => "kafka.juwatech.cn:9092"
    topics => ["juwatech-logs"]
    codec => "json"
  }
}

filter {
  # 提取嵌套字段
  json {
    source => "message"
    target => "log_data"
  }
  
  # 设置索引名按服务+日期
  mutate {
    add_field => { "[@metadata][index]" => "juwatech-%{[service]}-%{+YYYY.MM.dd}" }
  }
}

output {
  elasticsearch {
    hosts => ["http://es.juwatech.cn:9200"]
    index => "%{[@metadata][index]}"
    user => "elastic"
    password => "${ES_PASSWORD}"
  }
}

四、Kibana 查询实战

场景1:查找某用户返利失败原因

Kibana Discover 中输入:

userId: 10086 AND event: "REBATE_FAILED"

点击结果中的 traceId,使用 Trace ID 聚合视图 查看全链路日志。

场景2:监控佣金计算耗时异常

创建可视化图表,筛选:

event: "COMMISSION_CALCULATED" AND durationMs > 1000

设置告警规则,当 5 分钟内出现 10 次即触发企业微信通知。

场景3:统计各服务错误率

使用 Lens 可视化,按 service 分组,计算 level: ERROR 日志占比。

五、敏感信息脱敏

通过 Logback 自定义 converter 避免日志泄露:

package juwatech.cn.common.logging;

import ch.qos.logback.classic.pattern.ClassicConverter;
import ch.qos.logback.classic.spi.ILoggingEvent;

public class MaskingConverter extends ClassicConverter {
    @Override
    public String convert(ILoggingEvent event) {
        String msg = event.getFormattedMessage();
        // 脱敏手机号、身份证等
        return msg.replaceAll("(13[0-9]|14[01456879]|15[0-35-9]|16[2567]|17[0-8]|18[0-9]|19[0-35-9])\\d{8}",
            "$1****");
    }
}

在 logback.xml 中引用:

<conversionRule conversionWord="maskedMsg" converterClass="juwatech.cn.common.logging.MaskingConverter"/>
<encoder>
  <pattern>%maskedMsg</pattern>
</encoder>

六、日志采样与存储优化

  • 高频 INFO 日志按 10% 采样写入 ES,其余仅存本地文件;
  • 错误日志 100% 采集;
  • ES 索引生命周期管理(ILM):热数据保留 7 天,温数据转冷存储 30 天后删除。

该体系上线后,平均故障定位时间从 45 分钟缩短至 90 秒,运维效率显著提升。

本文著作权归 微赚淘客系统3.0 研发团队,转载请注明出处!

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐