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

如何在Scala自定义Log4j Appender中实现Elasticsearch批量日志上报

实现方案

完全可以按指定批次大小聚合日志后统一调用ES Bulk API提交,这也是生产环境常用的优化手段,可以大幅降低HTTP请求次数,提升上报性能。


核心实现思路

  • 在Appender内部维护线程安全的缓冲区,每次append调用时仅将构造好的单条bulk报文写入缓冲区,不直接发请求
  • 每次写入缓冲区后判断当前缓冲条数是否达到预设的批次阈值,达到就触发批量提交、清空缓冲区
  • 额外配置超时刷新逻辑,低流量场景下就算批次没满,到指定时间也主动提交,避免日志长时间滞留
  • 重写close方法,Appender销毁时强制提交缓冲区残留的所有日志,避免丢数据

优化后代码示例

import org.apache.log4j.AppenderSkeleton
import org.apache.log4j.spi.LoggingEvent
import java.util.concurrent.CopyOnWriteArrayList
import java.util.concurrent.atomic.AtomicLong
import java.util.Timer
import java.util.TimerTask

class KibanaAppender extends AppenderSkeleton {
  // 可配置参数,也可以通过Log4j配置文件注入
  private val batchSize: Int = 100 // 批次大小,达到就提交
  private val flushIntervalMs: Long = 5000 // 最长滞留时间5秒,到点就提交
  private val indexName: String = "your_log_index"

  // 线程安全的缓冲区
  private val buffer = new CopyOnWriteArrayList[String]()
  private val lastFlushTime = new AtomicLong(System.currentTimeMillis())
  private val lock = new Object()

  // 初始化定时刷新任务
  private val flushTimer = new Timer("KibanaAppender-Flush-Thread", true)
  flushTimer.scheduleAtFixedRate(new TimerTask {
    override def run(): Unit = {
      if (System.currentTimeMillis() - lastFlushTime.get() >= flushIntervalMs && buffer.size() > 0) {
        flush()
      }
    }
  }, flushIntervalMs, flushIntervalMs)

  override def append(event: LoggingEvent): Unit = {
    // 构造单条bulk报文
    val message = event.getMessage.toString
    val level = event.getLevel.toString
    val timeStamp = new java.text.SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSSZ").format(event.getTimeStamp)
    val singleDoc =
      s"""|{"create": {"_index" : "$indexName"}}
          |{"message":"$message","level":"$level","@timestamp":"$timeStamp"}
          |""".stripMargin

    buffer.add(singleDoc)
    // 达到批次阈值触发提交
    if (buffer.size() >= batchSize) {
      flush()
    }
  }

  // 批量提交方法
  private def flush(): Unit = {
    lock.synchronized {
      if (buffer.isEmpty) return
      // 拼接完整bulk请求,ES要求bulk报文最后必须有空行
      val bulkContent = buffer.toArray.mkString("\n") + "\n"
      try {
        HttpService.post(json = bulkContent)
        lastFlushTime.set(System.currentTimeMillis())
      } catch {
        case e: Exception =>
          // 这里可以加失败重试、或者落本地文件兜底逻辑,避免日志丢失
          e.printStackTrace()
      } finally {
        buffer.clear()
      }
    }
  }

  override def close(): Unit = {
    if (!this.closed) {
      // 关闭定时器,提交残留日志
      flushTimer.cancel()
      flush()
      this.closed = true
    }
  }

  override def requiresLayout(): Boolean = true
}

注意事项

  • 代码里的batchSize和flushIntervalMs可以根据你的业务流量调整,也可以改成通过Log4j的配置项动态注入,不需要硬编码
  • 生产环境建议加提交失败的兜底逻辑,比如重试2-3次还是失败的话,把当前批次的日志写入本地临时文件,后续可以人工补导
  • 如果日志量特别大,可以考虑把缓冲区改成有界队列,避免服务宕机时缓冲区里的日志丢失过多

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 04:30:04