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

Go-GST中AppSrc替代UDPSrc接收RTP流报错:不支持非TIME格式Segment

GStreamer AppSrc替代UDPSrc时Segment格式错误的解决方法

问题场景

我有如下GStreamer发送管道:

gst-launch-1.0 audiotestsrc ! audioconvert ! audioresample ! opusenc ! rtpopuspay ! udpsink host=127.0.0.1 port=5000

对应的接收管道运行正常:

gst-launch-1.0 udpsrc caps=application/x-rtp port=5000 ! rtpopusdepay ! opusdec ! autoaudiosink

使用go-gst在Golang中实现接收逻辑时,基于UDPSrc的版本可以正常运行,但改用Golang UDP监听器监听:5000,并通过AppSrc手动推送字节数据时,出现Segment with non-TIME format not supported错误。

原始实现代码可切换runAppSrc()(报错)和runUdpSrc()(正常):

package main

import (
    "fmt"
    "log"
    "net"
    "time"

    "github.com/go-gst/go-glib/glib"
    "github.com/go-gst/go-gst/gst"
    "github.com/go-gst/go-gst/gst/app"
)

func createAppSrcPipeline() (*gst.Pipeline, *app.Source, error) {
    gst.Init(nil)

    // Create a pipeline
    pipeline, err := gst.NewPipeline("")
    if err != nil {
        return nil, nil, err
    }

    // Create the elements
    elems, err := gst.NewElementMany("appsrc", "rtpopusdepay", "opusdec", "autoaudiosink")
    if err != nil {
        return nil, nil, err
    }

    // Add the elements to the pipeline and link them
    err = pipeline.AddMany(elems...)
    if err != nil {
        return nil, nil, err
    }
    err = gst.ElementLinkMany(elems...)
    if err != nil {
        return nil, nil, err
    }

    caps := gst.NewEmptySimpleCaps("application/x-rtp")
    caps.SetValue("media", "audio")
    caps.SetValue("clock-rate", 48000)
    caps.SetValue("payload", 96)
    caps.SetValue("encoding-name", "OPUS")

    src := app.SrcFromElement(elems[0])
    src.SetFormat(gst.FormatTime)
    src.SetDoTimestamp(true)
    src.SetLive(true)
    src.SetCaps(caps)

    return pipeline, src, nil
}

func createPipeline() (*gst.Pipeline, error) {
    gst.Init(nil)

    // Create a pipeline
    pipeline, err := gst.NewPipeline("")
    if err != nil {
        return nil, err
    }

    // Create the elements
    elems, err := gst.NewElementMany("udpsrc", "rtpopusdepay", "opusdec", "autoaudiosink")
    if err != nil {
        return nil, err
    }

    // Add the elements to the pipeline and link them
    err = pipeline.AddMany(elems...)
    if err != nil {
        return nil, err
    }
    err = gst.ElementLinkMany(elems...)
    if err != nil {
        return nil, err
    }

    caps := gst.NewEmptySimpleCaps("application/x-rtp")

    src := elems[0]
    src.Set("caps", caps)
    src.Set("port", 5000)

    return pipeline, nil
}

func handleMessage(msg *gst.Message) error {
    switch msg.Type() {
    case gst.MessageEOS:
        return app.ErrEOS
    case gst.MessageError:
        gerr := msg.ParseError()
        if debug := gerr.DebugString(); debug != "" {
            fmt.Println(debug)
        }
        return gerr
    }
    return nil
}

func mainLoop(loop *glib.MainLoop, pipeline *gst.Pipeline) error {
    // Start the pipeline

    // Due to recent changes in the bindings - the finalizers might fire on the pipeline
    // prematurely when it's passed between scopes. So when you do this, it is safer to
    // take a reference that you dispose of when you are done. There is an alternative
    // to this method in other examples.
    pipeline.Ref()
    defer pipeline.Unref()

    pipeline.SetState(gst.StatePlaying)

    // Retrieve the bus from the pipeline and add a watch function
    pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool {
        if err := handleMessage(msg); err != nil {
            fmt.Println(err)
            loop.Quit()
            return false
        }
        return true
    })

    loop.Run()

    return nil
}

func main() {
    runAppSrc()
    //runUdpSrc()
}

func runUdpSrc() {
    pipeline, err := createPipeline()
    if err != nil {
        return
    }

    ml := glib.NewMainLoop(glib.MainContextDefault(), false)

    if err = mainLoop(ml, pipeline); err != nil {
        fmt.Println("ERROR!", err)
    }
}

func runAppSrc() {
    pipeline, src, err := createAppSrcPipeline()
    if err != nil {
        log.Fatal(err)
    }

    go func() {
        ml := glib.NewMainLoop(glib.MainContextDefault(), false)
        if err := mainLoop(ml, pipeline); err != nil {
            fmt.Println("ERROR!", err)
        }
    }()

    // 等待管道初始化完成,发送TIME格式的Segment事件
    time.Sleep(100 * time.Millisecond)
    segment := gst.NewSegment()
    segment.Init(gst.FormatTime, 0, gst.CLOCK_TIME_NONE, 0)
    src.PushEvent(gst.NewEventSegment(segment))

    // listen to incoming udp packets
    pc, err := net.ListenPacket("udp", ":5000")
    if err != nil {
        log.Fatal(err)
    }
    defer pc.Close()

    for {
        buf := make([]byte, 1400)
        n, _, err := pc.ReadFrom(buf)
        if err != nil {
            continue
        }

        buffer := gst.NewBufferWithSize(int64(n))
        // 修正时间戳单位为纳秒
        buffer.SetPresentationTimestamp(gst.ClockTime(time.Now().UnixNano()))
        buffer.Map(gst.MapWrite).WriteData(buf[:n])
        buffer.Unmap()

        flow := src.PushBuffer(buffer)
        if flow == gst.FlowError {
            fmt.Println("Push buffer error:", flow)
            break
        }
    }
}

设置GST_DEBUG=*:2时的错误输出:

0:00:01.014558389 180022 0x7f404c000b90 ERROR       rtpbasedepayload gstrtpbasedepayload.c:970:gst_rtp_base_depayload_handle_event:<rtpopusdepay0> Segment with non-TIME format not supported
0:00:01.014574831 180022 0x7f404c000b90 ERROR       rtpbasedepayload gstrtpbasedepayload.c:970:gst_rtp_base_depayload_handle_event:<rtpopusdepay0> Segment with non-TIME format not supported
0:00:01.014582021 180022 0x7f404c000b90 WARN                 basesrc gstbasesrc.c:3132:gst_base_src_loop:<appsrc0> error: Internal data stream error.
0:00:01.014585366 180022 0x7f404c000b90 WARN                 basesrc gstbasesrc.c:3132:gst_base_src_loop:<appsrc0> error: streaming stopped, reason error (-5)
0:00:01.014607431 180022 0x7f404c000b90 ERROR       rtpbasedepayload gstrtpbasedepayload.c:970:gst_rtp_base_depayload_handle_event:<rtpopusdepay0> Segment with non-TIME format not supported
0:00:01.014609835 180022 0x7f404c000b90 ERROR       rtpbasedepayload gstrtpbasedepayload.c:970:gst_rtp_base_depayload_handle_event:<rtpopusdepay0> Segment with non-TIME format not supported
0:00:01.014611828 180022 0x7f404c000b90 ERROR       rtpbasedepayload gstrtpbasedepayload.c:970:gst_rtp_base_depayload_handle_event:<rtpopusdepay0> Segment with non-TIME format not supported
0:00:01.014631850 180022 0x7f404c000b90 WARN           audiobasesink gstaudiobasesink.c:1117:gst_audio_base_sink_wait_event:<autoaudiosink0-actual-sink-pulse> error: Sink not negotiated before eos event.

错误原因

rtpopusdepay要求流的Segment事件必须使用TIME格式,但AppSrc默认不会自动发送符合要求的Segment事件,且代码中存在两个关键问题:

  1. 时间戳单位错误:用毫秒设置GStreamer的ClockTime(实际单位是纳秒),导致时间戳无效。
  2. 未手动初始化Segment事件:UDPSrc会自动发送TIME格式的Segment事件初始化流,而AppSrc需要手动触发该事件。

修复方案

  1. 修正时间戳单位:将time.Now().UnixMilli()改为time.Now().UnixNano(),符合GStreamer ClockTime的纳秒要求。
  2. 发送TIME格式的Segment事件:在管道启动后,给AppSrc推送一个初始化的Segment事件,指定格式为TIME。
  3. 优化错误处理:避免外部err变量的作用域问题,在goroutine内部处理mainLoop的错误。

验证方法

运行修复后的Golang程序,同时启动发送管道,音频可以正常播放,不再出现Segment格式错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 16:57:08