You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为Karafka Server中的Rails日志实现跟踪ID?

为Karafka添加消费跟踪ID方案

1. 实现Karafka消费中间件生成跟踪ID

创建一个中间件,在消费批次启动时生成UUID作为跟踪ID存入线程上下文,消费完成后清理变量避免内存泄漏:

# app/middlewares/karafka_trace_id_middleware.rb
class KarafkaTraceIdMiddleware
  def call_consumer
    # 生成UUID作为跟踪ID
    Thread.current[:karafka_trace_id] = SecureRandom.uuid
    # 执行原消费逻辑
    yield
  ensure
    # 消费结束后清理线程变量
    Thread.current[:karafka_trace_id] = nil
  end
end

在karafka.rb中注册该中间件:

class KarafkaApp < Karafka::App
  kafka_config = Settings.kafka
  max_payload_size = 7_000_000 # 设置最大 payload 为7MB

  setup do |config|
    config.kafka = {
      'bootstrap.servers': ENV['KAFKA_BROKERS_SCRAM'],
      'max.poll.interval.ms': 1200000
    }
    config.client_id = kafka_config['client_id']
    config.concurrency = 4
    config.consumer_persistence = !Rails.env.development?

    config.producer = ::WaterDrop::Producer.new do |producer_config|
      producer_config.kafka = ::Karafka::Setup::AttributesMap.producer(config.kafka.dup)
      producer_config.max_payload_size = max_payload_size
      producer_config.kafka[:'message.max.bytes'] = max_payload_size
    end

    # 注册跟踪ID中间件
    config.middleware.append KarafkaTraceIdMiddleware
  end

  routes.draw do
    topic some_topic_1.to_sym do
      consumer SomeTopic1Consumer
    end

    # ... 其他路由 ...
  end
end

2. 修改Rails日志配置,统一跟踪ID输出

调整Rails日志格式,让它自动识别HTTP请求的request_id和Karafka消费的karafka_trace_id,实现统一的日志输出格式:

# config/application.rb
module YourAppName
  class Application < Rails::Application
    # ... 其他配置 ...

    # 自定义日志格式,优先使用Karafka跟踪ID,无则使用request_id
    config.log_formatter = proc do |severity, timestamp, progname, msg|
      trace_id = Thread.current[:karafka_trace_id] || Thread.current[:request_id]
      "[#{timestamp}] #{severity} [trace_id: #{trace_id}] #{msg}\n"
    end

    # 若使用TaggedLogging,可替换成以下配置
    # config.log_tags = [
    #   proc { Thread.current[:karafka_trace_id] || :request_id }
    # ]
  end
end

3. 可选:实现跟踪ID跨消息链路传递

如果需要在消息生产-消费的链路中传递跟踪ID,可修改中间件优先读取消息头的trace_id,没有再生成新ID:

# app/middlewares/karafka_trace_id_middleware.rb
class KarafkaTraceIdMiddleware
  def call_consumer
    # 优先从消息头获取trace_id,无则生成新ID
    batch = Karafka::Processing::Current.batch
    trace_id = batch.messages.first.headers['trace_id'] || SecureRandom.uuid
    Thread.current[:karafka_trace_id] = trace_id
    yield
  ensure
    Thread.current[:karafka_trace_id] = nil
  end
end

在消费者发送消息时,将当前跟踪ID写入消息头:

# 消费者中的生产逻辑示例
producer.produce_async(
  topic: 'target_topic',
  payload: data.to_json,
  headers: { 'trace_id' => Thread.current[:karafka_trace_id] }
)

效果验证

启动Karafka服务器消费消息后,日志会输出带跟踪ID的内容,格式与HTTP请求日志保持一致:

[2024-05-20 14:30:00] INFO [trace_id: 550e8400-e29b-41d4-a716-446655440000] 消息处理完成

内容的提问来源于stack exchange,提问作者user23353243

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 18:07:30