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

如何在Fiber框架中实时读取HTTP流式请求体数据?

问题根因

c.Body() 方法默认会等待整个HTTP请求体完全接收完成、全部缓存到内存后才会返回内容。ffmpeg推的是长流式请求,只要ffmpeg进程不终止、请求不结束,这个方法就会一直阻塞,自然拿不到实时传输的数据。

实现方案

Fiber底层基于fasthttp,原生支持请求体流式读取,核心是不要触发全量请求体缓存逻辑,直接读取底层TCP流的数据边收边处理。

前置配置

初始化Fiber实例时必须开启流式请求体相关配置,禁用默认的预解析、全量缓存逻辑:

app := fiber.New(fiber.Config{
    BodyLimit: 0, // 移除默认请求体大小限制,适配长流场景
    StreamRequestBody: true, // 核心配置:开启请求体流式读取
    DisablePreParseMultipartForm: true, // 禁用表单预解析,避免触发全量缓存
})

注意事项

  • 推流路由的处理逻辑中,绝对不能调用c.Body()、c.BodyParser()、c.FormValue()这类会触发全量请求体读取的方法,一旦调用就会阻塞到流结束。
  • 直接通过c.Context().RequestBodyStream()拿到请求体的可读流,循环读取字节块即可实时拿到ffmpeg推送的数据,读到io.EOF就代表推流结束。

完整示例代码

package main

import (
	"io"
	"log"

	"github.com/gofiber/fiber/v2"
	"github.com/gofiber/websocket/v2"
)

func main() {
	app := fiber.New(fiber.Config{
		BodyLimit:                    0,
		StreamRequestBody:            true,
		DisablePreParseMultipartForm: true,
	})

	// 拉流端websocket路由
	app.Get("/ws", websocket.New(func(c *websocket.Conn) {
		defer c.Close()
		// 生产环境此处实现连接鉴权、流标识绑定、加入对应流的转发连接池逻辑
		for {
			_, _, err := c.ReadMessage()
			if err != nil {
				break
			}
		}
	}))

	// ffmpeg推流入口
	app.Post("/push", func(c *fiber.Ctx) error {
		c.Status(200)
		// 获取底层请求体流
		bodyStream := c.Context().RequestBodyStream()
		// 读缓冲区,大小可根据码率调整,一般32KB~128KB即可
		buf := make([]byte, 32*1024)

		// 示例简化:生产环境此处从连接池取对应流的所有websocket连接做广播
		var targetWsConn *websocket.Conn

		for {
			n, readErr := bodyStream.Read(buf)
			// 先处理读到的有效数据
			if n > 0 {
				chunk := buf[:n]
				// 实时将数据块转发到websocket
				if targetWsConn != nil {
					if err := targetWsConn.WriteMessage(websocket.BinaryMessage, chunk); err != nil {
						log.Printf("websocket write failed: %v", err)
						break
					}
				}
				// 可在此扩展其他逻辑:切片录制、转推其他平台、实时转码等
			}
			// 处理流读取错误/结束
			if readErr != nil {
				if readErr != io.EOF {
					log.Printf("stream read failed: %v", readErr)
				} else {
					log.Println("ffmpeg push stream ended")
				}
				break
			}
		}
		return nil
	})

	log.Fatal(app.Listen(":3000"))
}

优化建议

  • ffmpeg推流时添加低延迟参数,避免ffmpeg自身缓存导致的延迟:
    ffmpeg -re -i <你的输入源> -c copy -f flv -fflags nobuffer -flags low_delay http://127.0.0.1:3000/push
  • 生产环境需要实现流ID与websocket连接池的映射,支持多流、多拉流端的广播逻辑,用channel做异步转发,避免慢拉流端阻塞推流读取逻辑。
  • 增加超时、断连检测、资源回收逻辑,避免长连接场景下的内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 11:36:16