如何为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
相关产品推荐
相关产品推荐

