很多人刚开始搭服务器的时候,排查问题的方式是这样的:
- SSH 登上服务器
tail -f /var/log/app.log- 盯着屏幕等报错
- 出问题了,赶紧复制粘贴,发给同事
要是服务器有三台呢?要是凌晨三点你在被窝里突然告警呢?
其实这个问题在业界早就有成熟的解法——日志采集 Agent。它是一个跑在你服务器上的轻量小程序,负责实时读取日志文件,然后把日志推送到一个中央平台,你只需要打开浏览器就能看到所有服务器的日志。
这篇文章就来教你从零开始,用 Go 语言实现一个这样的 Agent。最终产物是一个单文件二进制程序,编译后约 6MB,扔到任何 Linux 服务器上就能跑,不需要安装任何依赖。
一、先把问题想清楚,再动手写代码
很多人喜欢上来就开始写代码,结果写到一半发现设计有问题,推倒重来。所以我们先花几分钟把整个 Agent 需要解决的问题梳理一遍。
Agent 需要做哪些事?
① 实时监听日志文件,感知新内容
日志文件是不断追加内容的,Agent 需要一直盯着文件,有新的一行就立刻读出来。听起来简单,但有个细节:如果每次都从文件头开始读,那历史日志会被重复发送。所以 Agent 启动的时候,应该直接跳到文件末尾,只采集"我启动之后新产生的日志"。
② 批量上报,而不是来一条发一条
假设你的服务器在高峰期每秒产生 1000 行日志,如果每行都发一次 HTTP 请求,那就是每秒 1000 次网络请求。别说平台扛不住,你自己服务器的网络都会被打爆。
正确的做法是把日志攒一批再一起发。但是攒多少?攒多久?这里有个小技巧,后面详细讲。
③ 网络断了,日志不能丢
这是很多人没想到的场景。如果 Agent 往平台发日志的时候,刚好网络抖动了一下,这批日志就会发送失败。直接丢掉显然不行,我们需要失败重试机制。
好,三个核心问题清楚了。我们把整个数据流画出来:
日志文件 (app.log)
│
│ 一直盯着文件,有新行就读出来
▼
┌─────────────┐
│ Tailer │ 负责"监听"
│ (采集层) │
└──────┬──────┘
│ 把每一行写入一个内存队列(channel)
▼
┌─────────────┐
│ Batcher │ 负责"攒批"
│ (缓冲层) │ 积攒 100 条 OR 超过 500ms
└──────┬──────┘
│ 打包成一个 JSON 请求体
▼
┌─────────────┐
│ Sender │ 负责"发送"
│ (上报层) │ HTTP POST + 失败重试
└──────┬──────┘
│
▼
平台的接收接口
POST /api/v1/ingest
三层结构,职责完全分离。Tailer 不关心怎么发送,Sender 不关心怎么读文件,Batcher 在中间负责协调节奏。这样的设计,每一层都可以独立修改和优化,后续扩展也很方便。
二、初始化项目
Agent 是一个独立的 Go 程序,和你的后端平台完全分开,有自己的 main 函数和目录结构。
bash
# 创建目录结构
mkdir -p agent/tailer
mkdir -p agent/sender
# 初始化 go module(agent 是独立编译的,有自己的 module)
cd agent
go mod init logstream-agent
最终的目录结构长这样:
agent/
├── go.mod
├── go.sum
├── main.go # 程序入口:解析参数、串联三层
├── tailer/
│ └── file.go # Tailer:文件监听
└── sender/
└── http.go # Sender:HTTP 上报
三、定义数据结构
在开始写三层逻辑之前,先定义一个贯穿全局的数据结构——LogEntry,它代表从文件里采集到的一条日志。
go
// agent/main.go
package main
import "time"
// LogEntry 代表一条采集到的日志
type LogEntry struct {
Message string `json:"message"` // 日志的原文内容
Source string `json:"source"` // 来自哪个文件,例如 /var/log/app.log
Timestamp time.Time `json:"timestamp"` // 采集时间
}
这里有一个值得解释的设计决策:时间戳用服务端(Agent 本机)的 time.Now(),而不是日志文件里的时间。
为什么?
日志文件里的时间格式五花八门:有的写 2026-07-09 00:26:59,有的写 Jul 9 00:26:59,有的甚至没有时间,只有一段文字。解析这些格式需要写大量的适配代码,而且并不可靠。
用 Agent 本机的 time.Now() 虽然不是日志"产生"的精确时间,但它是日志被"采集"的时间,在大多数业务场景下已经足够精准(误差在几百毫秒之内)。
四、Tailer:实现文件实时监听
原理:tail -f 是怎么工作的?
你一定用过 tail -f app.log 这个命令。它的原理其实非常简单:
- 打开文件,把读取指针移到文件末尾
- 不断尝试读取新内容
- 如果没有新内容,等一小会儿(比如 100ms),再继续尝试
就是个死循环里的轮询,没有什么神奇之处。我们用 Go 来复现这个行为:
go
// agent/tailer/file.go
package tailer
import (
"bufio"
"fmt"
"os"
"time"
)
// Tail 持续监听 filePath 文件,将每一行新日志写入 out channel
// 这个函数会一直阻塞运行,直到发生不可恢复的错误
func Tail(filePath string, out chan<- string) error {
f, err := os.Open(filePath)
if err != nil {
return fmt.Errorf("打开文件失败: %w", err)
}
defer f.Close()
// 重点:把文件的读取指针移到末尾
// os.SEEK_END 表示"从文件末尾开始",偏移量 0 表示"就在末尾"
// 这样 Agent 启动后,只会读取新产生的日志,不会把历史日志全部重发
if _, err := f.Seek(0, os.SEEK_END); err != nil {
return fmt.Errorf("定位文件末尾失败: %w", err)
}
scanner := bufio.NewScanner(f)
for {
if scanner.Scan() {
// 读到了新的一行
line := scanner.Text()
// 跳过空行(日志文件里偶尔会有空行)
if line != "" {
out <- line // 把这行发送到 channel,等待 Batcher 处理
}
} else {
// scanner.Scan() 返回 false,说明当前没有新内容
// 休息 100ms 后继续轮询,避免 CPU 空转
time.Sleep(100 * time.Millisecond)
}
}
}
几个细节说明:
为什么用 bufio.Scanner 而不是 ioutil.ReadAll?
ioutil.ReadAll 会把整个文件内容一次性读入内存,对于一直在增长的日志文件来说,这会占用大量内存,而且没法做到"读一行发一行"。bufio.Scanner 是按行读取的,内存占用极低,非常适合这个场景。
chan<- string 是什么意思?
这是 Go 的语法,表示这是一个只写的 channel。在 Tail 函数里,我们只往 channel 里发数据(out <- line),不从里面读。用 chan<- string 而不是 chan string,是为了让编译器帮你检查:防止你不小心在这个函数里反向读取 channel,出现逻辑错误。
处理日志轮转(Log Rotation)
上面的代码在大多数情况下都能正常工作,但有一个生产环境里很常见的问题:日志轮转(Log Rotation)。
很多服务器会配置 logrotate 工具,它会在每天凌晨把 app.log 重命名为 app.log.2026-07-09,然后创建一个新的空 app.log 文件,供程序继续写入。
如果 Agent 还在用旧的文件描述符读取已经被重命名的那个文件,它会发现"永远没有新内容了",但实际上新的日志已经在写入新的 app.log 了。
解决方案是检测文件的 inode 是否发生变化。在 Linux 文件系统中,每个文件都有一个唯一的编号叫 inode。当 app.log 被重命名为 app.log.old,然后创建新的 app.log 时,新文件的 inode 会和旧文件不同。
我们可以定期检测这个变化:
go
// agent/tailer/file.go(完整版,支持日志轮转检测)
package tailer
import (
"bufio"
"fmt"
"os"
"syscall"
"time"
)
// getInode 获取文件的 inode 编号(Linux 文件系统的唯一标识)
func getInode(f *os.File) uint64 {
info, err := f.Stat()
if err != nil {
return 0
}
// 将 os.FileInfo 转换为底层的 syscall.Stat_t 结构,里面有 inode 信息
stat, ok := info.Sys().(*syscall.Stat_t)
if !ok {
return 0
}
return stat.Ino
}
func Tail(filePath string, out chan<- string) error {
f, err := os.Open(filePath)
if err != nil {
return fmt.Errorf("打开文件失败: %w", err)
}
defer f.Close()
if _, err := f.Seek(0, os.SEEK_END); err != nil {
return fmt.Errorf("定位文件末尾失败: %w", err)
}
currentInode := getInode(f) // 记录当前文件的 inode
scanner := bufio.NewScanner(f)
checkTicker := time.NewTicker(3 * time.Second) // 每 3 秒检查一次是否发生了日志轮转
defer checkTicker.Stop()
for {
if scanner.Scan() {
line := scanner.Text()
if line != "" {
out <- line
}
} else {
// 没有新内容,趁这个空隙检查一下文件是否被轮转了
select {
case <-checkTicker.C:
// 尝试打开同名文件,看看它的 inode 是否和当前的一样
newF, err := os.Open(filePath)
if err != nil {
// 文件暂时不存在(轮转间隙),等一下再试
break
}
newInode := getInode(newF)
if newInode != currentInode {
// inode 不同,说明文件被轮转了!
fmt.Printf("[Tailer] 检测到日志轮转,切换到新文件: %s\n", filePath)
f.Close()
f = newF
currentInode = newInode
scanner = bufio.NewScanner(f)
// 新文件从头开始读(轮转后的新文件是空的,没有历史数据)
} else {
newF.Close()
}
default:
// 不到检查时间,先休眠
time.Sleep(100 * time.Millisecond)
}
}
}
}
五、Batcher:实现智能批量缓冲
Batcher 是整个 Agent 里最有意思的一层,它需要解决一个很经典的工程问题:
怎样在"实时性"和"吞吐量"之间找到平衡?
- 如果每来一行日志就立刻发送:实时性最好,但网络请求太频繁,高并发下会出问题。
- 如果积攒 10000 行再发送:吞吐量很好,但日志可能要等几分钟才能在平台上看到,失去了"实时"的意义。
业界通用的解法是双触发机制:
只要满足以下任意一个条件,立刻发送当前积攒的所有日志:
- 条件 A:积攒的日志条数达到了上限(比如 100 条)
- 条件 B:距离上次发送已经超过了一段时间(比如 500ms)
条件 A 保证高峰期的吞吐量,条件 B 保证低谷期的实时性。这个设计在 Kafka、Logstash 等成熟的日志系统中都有类似的应用。
用 Go 的 select + time.Ticker 来实现这个逻辑,非常优雅:
go
// agent/main.go 中的 batchAndSend 函数
func batchAndSend(server, token, source string, lines <-chan string) {
// 用来存储当前积攒的日志
var batch []LogEntry
// 定时器:每 500ms 触发一次
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
// flush 函数:把当前 batch 里的所有日志发送出去,然后清空 batch
flush := func() {
if len(batch) == 0 {
return // 如果 batch 是空的,什么都不做
}
// 发送这批日志(这里调用 Sender 层,后面会实现)
sender.SendWithRetry(server, token, batch, 3)
// 清空 batch,注意这里不是 batch = nil,而是 batch = batch[:0]
// 后者可以复用底层数组的内存,避免 GC 频繁分配和回收内存
batch = batch[:0]
}
// 主循环:通过 select 同时监听两个事件
for {
select {
case line := <-lines:
// 事件一:Tailer 发来了一行新日志
batch = append(batch, LogEntry{
Message: line,
Source: source,
Timestamp: time.Now(),
})
// 检查是否触发"数量上限"条件
if len(batch) >= 100 {
flush()
}
case <-ticker.C:
// 事件二:定时器到了 500ms
// 不管 batch 里有没有数据,尝试 flush 一次
flush()
}
}
}
这里有一个 Go 并发编程的小技巧值得说一下:select 是 Go 里处理多个 channel 事件最常用的模式。它会同时监听 case 里的所有 channel,哪个先有数据就执行哪个,类似于操作系统里的 epoll。
六、Sender:HTTP 上报与失败重试
6.1 基础发送逻辑
Sender 的核心任务很简单:把一批 LogEntry 打包成 JSON,通过 HTTP POST 发送到平台的接收接口。
go
// agent/sender/http.go
package sender
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"time"
)
// LogEntry 和 main.go 里的一样(实际开发中可以提取到公共包)
type LogEntry struct {
Message string `json:"message"`
Source string `json:"source"`
Timestamp time.Time `json:"timestamp"`
}
// 全局共用一个 http.Client,并设置超时时间
// 注意:http.Client 是线程安全的,可以复用
var httpClient = &http.Client{
Timeout: 10 * time.Second,
}
// Send 将一批日志发送到平台
func Send(serverAddr, token string, logs []LogEntry) error {
// 第一步:把日志列表序列化成 JSON
// 接口接受的格式是 {"logs": [...]}
body, err := json.Marshal(map[string]any{
"logs": logs,
})
if err != nil {
// json.Marshal 极少失败,除非数据里有不可序列化的类型
return fmt.Errorf("JSON 序列化失败: %w", err)
}
// 第二步:构建 HTTP POST 请求
req, err := http.NewRequest(
"POST",
serverAddr+"/api/v1/ingest",
bytes.NewReader(body), // 把 JSON 字节数组作为请求体
)
if err != nil {
return fmt.Errorf("构建 HTTP 请求失败: %w", err)
}
// 第三步:设置请求头
// X-Agent-Token 是身份验证凭证,平台通过它识别是哪个 Agent 在上报
req.Header.Set("X-Agent-Token", token)
req.Header.Set("Content-Type", "application/json")
// 第四步:发送请求
resp, err := httpClient.Do(req)
if err != nil {
// 网络故障(超时、连接拒绝等)
return fmt.Errorf("发送请求失败: %w", err)
}
defer resp.Body.Close()
// 第五步:检查响应状态码
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("平台返回异常状态码: %d", resp.StatusCode)
}
return nil
}
6.2 失败重试:为什么需要"退避"?
天真的重试策略是:失败了就立刻再试,还失败就再试……
这个策略有个严重的问题:如果服务器宕机了,大量 Agent 同时以极高的频率发起重试,会在服务器恢复的瞬间造成"雪崩"——服务器刚喘口气,就被一波请求又打倒了。
正确的做法是指数退避(Exponential Backoff):每次重试之间的等待时间指数级增长。
第 1 次失败 → 等待 1 秒后重试
第 2 次失败 → 等待 2 秒后重试
第 3 次失败 → 等待 4 秒后重试
第 4 次失败 → 等待 8 秒后重试
...
这样,既保证了短暂的网络抖动能快速恢复,又避免了长时间故障时的"重试风暴"。
go
// agent/sender/http.go(续)
// SendWithRetry 带指数退避的发送,最多重试 maxRetries 次
func SendWithRetry(serverAddr, token string, logs []LogEntry, maxRetries int) {
for attempt := 0; attempt < maxRetries; attempt++ {
err := Send(serverAddr, token, logs)
if err == nil {
// 发送成功,直接返回
return
}
// 计算下次重试的等待时间:1s, 2s, 4s, 8s...
// 用位移运算:1 << 0 = 1, 1 << 1 = 2, 1 << 2 = 4
waitTime := time.Duration(1<<uint(attempt)) * time.Second
fmt.Printf("[Sender] 第 %d 次发送失败,将在 %v 后重试。错误: %v\n",
attempt+1, waitTime, err)
time.Sleep(waitTime)
}
// 超过最大重试次数,这批日志只能丢弃了
// TODO: 生产环境中,这里可以把日志写入本地临时文件,网络恢复后补发
fmt.Printf("[Sender] 已重试 %d 次,放弃发送 %d 条日志\n", maxRetries, len(logs))
}
七、主程序:把三层串联起来
有了 Tailer、Batcher 和 Sender,主程序的工作就是解析命令行参数,然后把三层串起来。
go
// agent/main.go(完整版)
package main
import (
"flag"
"fmt"
"logstream-agent/sender"
"logstream-agent/tailer"
"os"
"time"
)
type LogEntry struct {
Message string `json:"message"`
Source string `json:"source"`
Timestamp time.Time `json:"timestamp"`
}
func main() {
// 定义命令行参数
serverAddr := flag.String("server", "", "平台地址,例如: http://127.0.0.1:8000")
token := flag.String("token", "", "Agent Token,在平台创建 Agent 时获取")
filePath := flag.String("file", "", "监听的日志文件路径,例如: /var/log/app.log")
flag.Parse()
// 校验必填参数
if *serverAddr == "" || *token == "" || *filePath == "" {
fmt.Println("缺少必要参数!")
fmt.Println()
fmt.Println("使用方式:")
fmt.Println(" logstream-agent -server <平台地址> -token <Agent Token> -file <日志文件路径>")
fmt.Println()
fmt.Println("示例:")
fmt.Println(" logstream-agent -server http://logs.example.com:8000 -token abc123 -file /var/log/app.log")
os.Exit(1)
}
fmt.Printf("[Agent] 启动成功\n")
fmt.Printf("[Agent] 监听文件: %s\n", *filePath)
fmt.Printf("[Agent] 上报目标: %s\n", *serverAddr)
// 创建一个带缓冲的 channel,作为 Tailer 和 Batcher 之间的"传送带"
// 缓冲大小设为 1000:当 Batcher 来不及处理时,Tailer 最多可以先积压 1000 行
// 超过 1000 行后,Tailer 会阻塞等待,而不是无限制地占用内存
lines := make(chan string, 1000)
// 在独立的 goroutine 中启动 Tailer
// Tailer 会一直阻塞运行,所以必须放在 goroutine 里
go func() {
if err := tailer.Tail(*filePath, lines); err != nil {
fmt.Printf("[Agent] Tailer 异常退出: %v\n", err)
os.Exit(1)
}
}()
// 主 goroutine 负责 Batcher(阻塞运行,程序不会退出)
batchAndSend(*serverAddr, *token, *filePath, lines)
}
func batchAndSend(server, token, source string, lines <-chan string) {
var batch []LogEntry
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
flush := func() {
if len(batch) == 0 {
return
}
// 注意:这里需要复制一份 batch 的内容,因为 SendWithRetry 是同步的
// 如果直接传 batch,在 SendWithRetry 运行期间,batch 可能被清空
toSend := make([]sender.LogEntry, len(batch))
for i, e := range batch {
toSend[i] = sender.LogEntry{
Message: e.Message,
Source: e.Source,
Timestamp: e.Timestamp,
}
}
sender.SendWithRetry(server, token, toSend, 3)
batch = batch[:0]
}
for {
select {
case line := <-lines:
batch = append(batch, LogEntry{
Message: line,
Source: source,
Timestamp: time.Now(),
})
if len(batch) >= 100 {
flush()
}
case <-ticker.C:
flush()
}
}
}
八、编译与部署
编译:交叉编译到 Linux
Go 最强大的特性之一就是交叉编译——你可以在 Mac 或 Windows 上,直接编译出能在 Linux 服务器上运行的二进制文件。
bash
# 在你的开发机上执行,编译出 Linux 64位 的可执行文件
GOOS=linux GOARCH=amd64 go build -o logstream-agent .
# 查看文件大小(Go 编译出的二进制包含所有依赖,无需额外安装)
ls -lh logstream-agent
# -rwxr-xr-x 1 user staff 6.1M logstream-agent
两个环境变量的含义:
GOOS=linux:目标操作系统是 LinuxGOARCH=amd64:目标 CPU 架构是 x86_64(绝大多数服务器的架构)
如果你的服务器是 ARM 架构(比如树莓派、部分云服务器),把 amd64 改成 arm64 即可。
部署:传到服务器,直接跑
bash
# 把编译好的文件传到服务器
scp logstream-agent root@your-server-ip:/usr/local/bin/
# SSH 到服务器
ssh root@your-server-ip
# 给文件赋予执行权限
chmod +x /usr/local/bin/logstream-agent
# 测试运行(前台跑,先看看有没有问题)
logstream-agent \
-server http://your-platform:8000 \
-token 你的AgentToken \
-file /var/log/nginx/access.log
生产部署:配置为 systemd 服务
直接在命令行跑有个问题——你一关 SSH 会话,Agent 就停了。生产环境里,我们应该把它配置成系统服务,开机自启、崩溃自动重启。
bash
# 创建 systemd 服务文件
cat > /etc/systemd/system/logstream-agent.service << 'EOF'
[Unit]
Description=LogStream Log Agent
# 等网络就绪后再启动
After=network.target
[Service]
Type=simple
# 启动命令
ExecStart=/usr/local/bin/logstream-agent \
-server http://your-platform:8000 \
-token 你的AgentToken \
-file /var/log/app.log
# 异常退出后,等 5 秒自动重启
Restart=on-failure
RestartSec=5s
# 以 root 用户运行(如果日志文件权限需要的话)
User=root
[Install]
WantedBy=multi-user.target
EOF
# 重新加载 systemd 配置
systemctl daemon-reload
# 设置开机自启
systemctl enable logstream-agent
# 立刻启动
systemctl start logstream-agent
# 查看运行状态
systemctl status logstream-agent
如果一切正常,你会看到类似这样的输出:
● logstream-agent.service - LogStream Log Agent
Loaded: loaded (/etc/systemd/system/logstream-agent.service; enabled)
Active: active (running) since Mon 2026-07-09 00:26:59 CST; 5s ago
Main PID: 12345 (logstream-agent)
九、测试验证
全部部署好之后,我们来验证一下整个链路是否通畅。
在服务器上模拟写入日志:
bash
# 每秒写一行日志到文件
while true; do
echo "$(date '+%Y-%m-%d %H:%M:%S') [INFO] 用户登录成功 user_id=12345" >> /var/log/app.log
echo "$(date '+%Y-%m-%d %H:%M:%S') [ERROR] 数据库连接超时 db=mysql" >> /var/log/app.log
sleep 1
done
打开你的日志平台 Dashboard,应该能看到这些日志以大约 500ms 的延迟实时出现在界面上。
十、总结与后续扩展
我们用大约 200 行 Go 代码,从零实现了一个具备以下能力的 Log Agent:
| 功能 | 实现方式 |
|---|---|
| 实时文件监听 | bufio.Scanner + 100ms 轮询 |
| 跳过历史日志 | 启动时 Seek 到文件末尾 |
| 日志轮转感知 | inode 变更检测 |
| 批量上报 | select + time.Ticker 双触发 |
| 失败重试 | 指数退避,最多 N 次 |
| 身份验证 | HTTP Header X-Agent-Token |
| 单文件部署 | Go 静态编译,零依赖 |
这只是一个 MVP 版本,后续还有很多可以扩展的方向:
- 📂 多文件监听:同时监听多个日志文件,例如 access.log 和 error.log
- 🐳 Docker 日志支持:通过 Docker API 采集容器的标准输出
- 💾 本地磁盘缓冲:网络长时间断开时,把日志暂存到本地文件,网络恢复后补发,实现真正的零丢失
- 🔍 日志级别识别:通过正则自动解析日志里的
[ERROR]、[WARN]、[INFO]标签 - 📊 上报指标:统计并定期打印"已采集行数 / 已发送批次 / 发送失败次数"等运行指标
希望这篇文章对你有帮助!如果有任何问题,欢迎在评论区讨论。