云计算百科
云计算领域专业知识百科平台

Monibuca 的 RTMP 特快专列

序:一条流媒体的丝绸之路

        在流媒体服务器的世界里,RTMP 就像一条古老的丝绸之路——数据从远方跋涉而来,穿越协议的重重关隘,最终抵达目的地。而今天,我们要跟随 Monibuca v5​ 的脚步,看看它是如何用 Go 语言打造这条"RTMP 特快专列"的。

第一:蓝图上的零件

        故事要从一个叫 package rtmp 的地方说起。这里住着几个核心角色:

角色

身份

职责

Client

全能通信兵

负责连接、握手、收发消息

Puller

拉流特工

从远程 RTMP 服务器拉取音视频流

Pusher

推流使者

将本地流转推到远程 RTMP 服务器

NetStream

底层引擎

处理 RTMP 块(chunk)和 AMF 编解码

        它们都藏在 m7s.live/v5 这个大工厂里,各司其职。

        解读源码链接

type Client struct {
NetStream
chunkSize int // 数据块大小,默认 4096 字节
u *url.URL // 解析后的 RTMP 地址
}

Client 是所有 RTMP 客户端的基础骨架。它继承了 NetStream 的通信能力,记住了数据块大小和服务器地址。

第二:URL 的拆解密码

        想象你是一个快递员,收到一个地址:rtmp://live.example.com/app/stream?token=abc。

        第一步,你得把它拆开:

func (c *Client) commonStart(ctx context.Context, addr string) (err error) {
c.u, err = url.Parse(addr)
// rtmp://live.example.com/app/stream?token=abc
// ↓ 拆解
// Host: live.example.com:1935
// Path: /app/stream
// AppName: app
// StreamName: stream?token=abc
}

代码会自动补全端口号:

  • rtmp:// → 默认 1935
  • rtmps:// → 默认 443(加密通道)

if strings.Count(c.u.Host, ":") == 0 {
if isRtmps {
c.u.Host += ":443"
} else {
c.u.Host += ":1935"
}
}

就像快递系统自动补全邮编一样,代码默默帮你把"不完整"的地址补齐了。

第三:三次握手,建立信任

        RTMP 协议有一个著名的"三次握手"(C0/C1/C2 + S0/S1/S2)。在代码中,这一步被封装在:

func (c *Client) commonRun(handler func(commander Commander) error) (err error) {
if err = c.ClientHandshake(); err != nil {
return
}
// 握手成功!发送块大小
err = c.SendMessage(RTMP_MSG_CHUNK_SIZE, Uint32Message(c.chunkSize))

        握手完成后,客户端会发送一个 connect​ 命令(AMF0 编码):

err = c.SendMessage(RTMP_MSG_AMF0_COMMAND, &CallMessage{
CommandMessage{"connect", 1},
map[string]any{
"app": c.AppName,
"flashVer": "monibuca/" + m7s.Version,
"swfUrl": c.u.String(),
"tcUrl": strings.TrimSuffix(c.u.String(), path) + "/" + c.AppName,
},
nil,
})

        这就像你在餐厅进门时说:"你好,我要连接 app 应用,我是 Monibuca v5 客户端。"服务器听到后会回一句 _result 或 _error。

第四:Puller 的拉流冒险 

        现在,故事的主角登场了——Puller(拉流特工)。

4.1 出发前的准备

func (p *Puller) Start() (err error) {
p.pullCtx.SetProgressStepsDefs(rtmpPullSteps)
addr := p.pullCtx.Connection.RemoteURL
err = p.pullCtx.Publish() // 向本地注册一个发布者

        Puller 先向本地系统"报到",说:"我要开始拉流了,请给我分配一个发布者身份。"

        然后它一步步推进进度条:

StepPublish → "Publishing stream"
StepURLParsing → "Parsing RTMP URL"
StepConnection → "Connecting to RTMP server"
StepHandshake → "Performing RTMP handshake"
StepStreaming → "Receiving media stream"

4.2 连接与握手

p.pullCtx.GoToStepConst(pkg.StepURLParsing)
err = p.commonStart(p.pullCtx.Context, addr)
p.pullCtx.GoToStepConst(pkg.StepConnection)
p.pullCtx.GoToStepConst(pkg.StepHandshake)
p.pullCtx.GoToStepConst(pkg.StepStreaming)

        每一步都像通关打卡,失败了就调用 p.pullCtx.Fail(err.Error()) 宣告任务失败。

4.3 进入 Run:发送 play 命令

func (p *Puller) Run() (err error) {
return p.commonRun(func(commander Commander) error {
switch response := commander.(type) {
case *ResponseCreateStreamMessage:
m := &PlayMessage{}
m.CommandMessage.CommandName = "play"
m.StreamName = ps[len(ps)-1] // 流名称
return p.SendMessage(RTMP_MSG_AMF0_COMMAND, m)
}
return nil
})
}

剧情高潮:当服务器返回 createStream 的响应后,Puller 立刻发送 play​ 命令。就像在视频网站点击"播放"按钮,数据流开始源源不断地涌来。

第五:Pusher 的推流使命 

        与 Puller 相反,Pusher(推流使者)​ 的任务是将本地流转推到远程服务器。

5.1 Start:轻装上阵

func (p *Pusher) Start() (err error) {
return p.commonStart(p.pushCtx.Context, p.pushCtx.Connection.RemoteURL)
}

        Pusher 的 Start 很简洁——只需要建立连接,剩下的交给 Run。

5.2 Run:createStream

case *ResponseCreateStreamMessage:
p.StreamID = response.StreamId
err = p.pushCtx.Subscribe() // 订阅本地流
return p.SendMessage(RTMP_MSG_AMF0_COMMAND, &PublishMessage{
CURDStreamMessage{…},
streamPath,
"live",
})

        先订阅本地流(拿到音视频数据),然后发送 publish​ 命令告诉远程服务器:"我要开始推流了!"

5.3 收到确认,开始传输

case *ResponsePublishMessage:
if response.Infomation["code"] == NetStream_Publish_Start {
p.Subscribe(p.pushCtx.Subscriber)
// ✅ 推流正式开始!
}

第六:工厂的流水线

        最后,看看这些角色是如何被"生产"出来的:

func NewPuller(_ config.Pull) m7s.IPuller {
ret := &Puller{
Client: Client{chunkSize: 4096},
}
ret.NetConnection = &NetConnection{}
ret.SetDescription(task.OwnerTypeKey, "RTMPPuller")
return ret
}

func NewPusher() m7s.IPusher {
ret := &Pusher{
Client: Client{chunkSize: 4096},
}
ret.NetConnection = &NetConnection{}
ret.SetDescription(task.OwnerTypeKey, "RTMPPusher")
return ret
}

        工厂函数 NewPuller 和 NewPusher 就像两条装配线,把 Client 骨架、NetConnection 引擎组装好,贴上标签,交付使用。

尾声:代码背后的哲学

        这段 RTMP 实现体现了几个设计思想:

  • 复用与分离:Client 封装了通用逻辑,Puller 和 Pusher 只需关注各自业务
  • 进度可观测:通过 StepDef 和 GoToStepConst,每一步都有清晰的进度反馈
  • 协议即状态机:commonRun 中的 switch 就是一个状态机,根据服务器响应推进状态
  • 错误处理即故事:每个 err != nil 都是一个"失败分支",让调用者知道哪里出了问题
  •         最后的话:RTMP 协议虽然"古老",但它依然是直播领域的中坚力量。Monibuca v5 用 Go 语言重新诠释了这条经典协议,让它在现代流媒体架构中继续发光发热。

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » Monibuca 的 RTMP 特快专列
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!