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

