FFmpeg com Scala: vídeo em tempo real
Use Scala para iniciar um fluxo RTSP de teste ao vivo, transformá-lo com FFmpeg e disponibilizar o fluxo HLS resultante por HTTP local. Você verá um padrão em movimento em tons de cinza com uma marca de tempo e manterá uma gravação de 12 segundos. O exemplo roda no Linux, não precisa de câmera e encerra os próprios serviços quando a demonstração termina.
O ponto central é a responsabilidade pelos processos: Scala inicia o FFmpeg e aguarda seu término
efetivo, mantendo um prazo limite e um hook de encerramento que podem interrompê-lo. Aguardar apenas
um Future do Scala não cancela um codificador externo. Trata-se de um único
fluxo somente de vídeo, em um único nível de qualidade, com o armazenamento em buffer comum do HLS;
ele não demonstra uma meta de latência de produção nem entrega com taxa de bits adaptativa.
Configurar Scala e FFmpeg
Este passo a passo foi testado no Linux x86-64 com OpenJDK 21.0.12.1, Scala CLI 1.17.1,
Scala 2.13.18, MediaMTX 1.21.1 e
FFmpeg/ffprobe/ffplay 9.0.1. Use Bash para os comandos. Instale uma
compilação do FFmpeg para Linux com libx264,
o multiplexador HLS e os filtros testsrc2, scale,
hue, drawbox e drawtext.
Você também precisa desse JDK, do cURL, do gzip, do tar, do awk, do
fc-match do fontconfig e de uma fonte TrueType instalada, como Liberation Sans.
O FFplay precisa de um ambiente de desktop gráfico para exibir a reprodução.
java -version && ffmpeg -version && ffprobe -version && ffplay -version &&
fc-match 'Liberation Sans'
Crie um novo projeto e baixe as versões fixadas do
Scala CLI e do
MediaMTX. O MediaMTX é o outro lado da conexão RTSP local:
o FFmpeg publica nele, e depois outro processo do FFmpeg lê esse fluxo.
Cole este bloco a partir de um diretório pai com permissão de escrita. Se
scala-live-demo já existir, a configuração é interrompida sem alterar esse diretório;
escolha outro diretório pai para recomeçar.
(
mkdir scala-live-demo &&
cd scala-live-demo &&
curl -fsSLo scala-cli.gz https://github.com/VirtusLab/scala-cli/releases/download/v1.17.1/scala-cli-x86_64-pc-linux.gz &&
gzip -d scala-cli.gz &&
chmod +x scala-cli &&
curl -fsSLo mediamtx.tar.gz https://github.com/bluenviron/mediamtx/releases/download/v1.21.1/mediamtx_v1.21.1_linux_amd64.tar.gz &&
tar -xzf mediamtx.tar.gz mediamtx &&
font_file=$(fc-match -f '%{file}' 'Liberation Sans') &&
test -r "$font_file" &&
cp "$font_file" font.ttf &&
./scala-cli version && ./mediamtx --version
)
Permaneça no diretório pai para os demais comandos. Cada bloco entre parênteses mantém locais as
mudanças de diretório e as opções do shell. Uma falha no download pode deixar o novo projeto
incompleto; inspecione-o e use outro diretório pai antes de tentar novamente. Salve o programa
completo abaixo como scala-live-demo/Live.scala. A primeira compilação baixa o compilador e as
bibliotecas do Scala; não há dependências adicionais da aplicação. Os comandos usam explicitamente
Live.scala como alvo; um projeto sbt que o englobe não é a entrada da compilação.
Integrar FFmpeg ao Scala
O programa usa o ProcessBuilder do Java com argumentos separados, em vez de
interpolar um comando de shell. Ele grava uma configuração privada do MediaMTX, vincula o HTTP à
interface de loopback e escolhe portas disponíveis por padrão. Somente a playlist e os nomes de
arquivo dos segmentos concluídos são disponibilizados; a configuração, as fontes e os logs não
ficam acessíveis por HTTP.
//> using scala 2.13.18
import com.sun.net.httpserver.{HttpExchange, HttpServer}
import java.net.{InetAddress, InetSocketAddress, ServerSocket}
import java.nio.file.{Files, LinkOption, Path, Paths}
import java.util.concurrent.TimeUnit.SECONDS
import scala.collection.mutable.ArrayBuffer
import scala.util.control.NonFatal
object Live {
private val children = ArrayBuffer.empty[(String, java.lang.Process)]
private var http: Option[HttpServer] = None
private var stopping = false
@volatile private var encoded = false
private def stopOwned(): Unit = synchronized {
stopping = true
http.foreach(_.stop(0))
http = None
for ((name, child) <- children.reverse if child.isAlive) {
child.destroy()
if (!child.waitFor(2, SECONDS)) {
child.destroyForcibly()
require(child.waitFor(2, SECONDS), s"Could not stop $name")
}
}
children.clear()
}
private def start(name: String, args: Seq[String], out: Path): java.lang.Process = synchronized {
require(!stopping, "Run is stopping")
val builder = new ProcessBuilder(args: _*)
.directory(out.toFile)
.redirectErrorStream(true)
.redirectOutput(out.resolve(s"$name.log").toFile)
// Ambient MediaMTX overrides must not reopen listeners disabled by our configuration.
builder.environment().keySet().removeIf(key => key.startsWith("MTX_"))
val child = builder.start()
children += name -> child
println(s"Started $name (PID ${child.pid()})")
child
}
private def serve(exchange: HttpExchange, out: Path): Unit = {
try {
val name = exchange.getRequestURI.getPath.stripPrefix("/")
val allowed = name == "stream.m3u8" || name.matches("segment-[0-9]{3}\\.ts")
val file = out.resolve(name)
if (exchange.getRequestMethod != "GET") exchange.sendResponseHeaders(405, -1)
else if (!allowed || !Files.isRegularFile(file)) exchange.sendResponseHeaders(404, -1)
else {
val contentType = if (name.endsWith(".ts")) "video/mp2t" else "application/vnd.apple.mpegurl"
exchange.getResponseHeaders.set("Content-Type", contentType)
exchange.getResponseHeaders.set("Cache-Control", "no-store")
val bytes = Files.readAllBytes(file)
exchange.sendResponseHeaders(200, bytes.length)
exchange.getResponseBody.write(bytes)
}
} finally exchange.close()
}
private def run(args: Array[String]): Unit = {
require(args.length >= 2 && args.length <= 4,
"Usage: Live <new-directory> <deadline-seconds> [rtsp-port] [http-port]")
val seconds = args(1).toInt
require(seconds >= 1 && seconds <= 120, "Deadline must be between 1 and 120 seconds")
val ports = args.drop(2).map(_.toInt)
require(ports.forall(port => port >= 0 && port <= 65535), "Invalid port")
val project = Paths.get(".").toRealPath()
val out = project.resolve(args(0)).normalize()
require(out.getParent == project, "Use a new direct child directory of this project")
require(!Files.exists(out, LinkOption.NOFOLLOW_LINKS), "Output directory already exists")
require(Files.isReadable(project.resolve("font.ttf")), "Missing readable font.ttf")
Files.createDirectory(out)
Files.copy(project.resolve("font.ttf"), out.resolve("font.ttf"))
val deadline = System.nanoTime() + SECONDS.toNanos(seconds)
def checkDeadline(): Unit =
require(System.nanoTime() < deadline, "Run deadline exceeded; output may be partial")
def alive(name: String, child: java.lang.Process): Unit =
require(child.isAlive, s"$name stopped (exit ${child.exitValue()}); inspect $name.log")
def waitUntil(ready: => Boolean)(check: => Unit): Unit = {
while (!ready) { checkDeadline(); check; Thread.sleep(100) }
checkDeadline()
}
val hook = new Thread(() => {
Console.err.println(if (encoded) "Stopped; completed HLS retained" else "Cancelled; output may be partial")
stopOwned()
})
Runtime.getRuntime.addShutdownHook(hook)
try {
println(s"Controller PID ${ProcessHandle.current().pid()}")
val requestedRtsp = ports.headOption.getOrElse(0)
val reservation = new ServerSocket(requestedRtsp, 1, InetAddress.getByName("127.0.0.1"))
val rtspPort = reservation.getLocalPort
reservation.close()
val server = synchronized {
require(!stopping, "Run is stopping")
val listener = HttpServer.create(new InetSocketAddress("127.0.0.1", ports.lift(1).getOrElse(0)), 0)
http = Some(listener)
listener.createContext("/", exchange => serve(exchange, out))
listener.start()
listener
}
val url = s"rtsp://127.0.0.1:$rtspPort/live"
val config = out.resolve("mediamtx.yml")
Files.writeString(config, s"""logLevel: info
rtsp: yes
rtspAddress: 127.0.0.1:$rtspPort
rtspTransports: [tcp]
rtmp: no
hls: no
webrtc: no
srt: no
moq: no
api: no
metrics: no
pprof: no
playback: no
paths:
live:
source: publisher
""")
val mtx = start("mediamtx", Seq(project.resolve("mediamtx").toString, config.toString), out)
def serverLog: String = Files.readString(out.resolve("mediamtx.log"))
waitUntil(serverLog.contains("[RTSP] started with listeners")) { alive("mediamtx", mtx) }
val common = Seq("ffmpeg", "-hide_banner", "-loglevel", "warning", "-nostdin")
val source = start("source", common ++ Seq(
"-re", "-f", "lavfi", "-i", "testsrc2=size=320x180:rate=10", "-t", "45",
"-an", "-c:v", "libx264", "-threads", "2", "-preset", "ultrafast",
"-tune", "zerolatency", "-g", "10", "-pix_fmt", "yuv420p",
"-f", "rtsp", "-rtsp_transport", "tcp", url), out)
waitUntil(serverLog.contains("is publishing to path 'live'")) {
alive("mediamtx", mtx); alive("source", source)
}
val filter = "scale=160:90,hue=s=0,drawbox=x=0:y=0:w=iw:h=22:color=black:t=fill," +
"drawtext=fontfile=font.ttf:text='Scala %{pts\\:hms}':fontcolor=white:fontsize=12:x=5:y=4"
val encoder = start("encoder", common ++ Seq(
"-n", "-rtsp_transport", "tcp", "-timeout", "3000000", "-i", url,
"-map", "0:v:0", "-an", "-vf", filter, "-filter_threads", "2", "-frames:v", "120",
"-c:v", "libx264", "-threads", "2", "-preset", "ultrafast", "-tune", "zerolatency",
"-pix_fmt", "yuv420p", "-g", "20", "-sc_threshold", "0", "-flags", "+cgop",
"-f", "hls", "-hls_time", "2", "-hls_playlist_type", "event",
"-hls_flags", "independent_segments+temp_file", "-hls_segment_filename", "segment-%03d.ts",
"-progress", "progress.txt", "stream.m3u8"), out)
val playlist = out.resolve("stream.m3u8")
waitUntil(Files.exists(playlist)) {
alive("mediamtx", mtx); alive("source", source); alive("encoder", encoder)
}
val playback = s"http://127.0.0.1:${server.getAddress.getPort}/stream.m3u8"
Files.writeString(out.resolve("playback.url"), playback + "\n")
println(s"Play: $playback")
waitUntil(!encoder.isAlive) { alive("mediamtx", mtx); alive("source", source) }
require(encoder.exitValue() == 0, "Encoder failed; inspect encoder.log")
val frames = Files.readString(out.resolve("progress.txt")).linesIterator
.filter(_.startsWith("frame=")).toVector.lastOption
require(frames.exists(_.stripPrefix("frame=").trim == "120"), "Incomplete video; output may be partial")
require(Files.readString(playlist).contains("#EXT-X-ENDLIST"), "Playlist was not finalized")
encoded = true
Files.writeString(out.resolve("complete.txt"), "120 frames encoded\n")
println("Encoding complete; HTTP remains available until the run deadline")
while (System.nanoTime() < deadline) Thread.sleep(100)
} finally {
stopOwned()
Runtime.getRuntime.removeShutdownHook(hook)
}
println("Completed; services stopped and HLS retained")
}
def main(args: Array[String]): Unit = {
try run(args)
catch { case NonFatal(error) => Console.err.println(s"Live demo failed: ${error.getMessage}"); sys.exit(1) }
}
}
Processamento de fluxos de vídeo em tempo real
O -re do FFmpeg regula o ritmo da fonte gerada, em vez de publicá-la na
velocidade máxima de codificação da CPU. O MediaMTX aceita a conexão do publicador em
/live; o codificador lê esse mesmo
caminho RTSP usando
transporte TCP. Este exemplo usa RTSP, em vez de pressupor que
RTMP ou uma câmera qualquer tenha o mesmo comportamento de entrada.
O filtro redimensiona para 160 por 90 pixels, remove as cores e desenha uma faixa preta com uma
marca de tempo decorrido. drawtext precisa de
suporte a fontes na sua compilação do FFmpeg, mesmo que o arquivo de fonte seja indicado
explicitamente. Compilações com fontconfig podem recorrer a outra fonte se esse arquivo for
inválido, então use a fonte instalada selecionada durante a configuração e confira a sobreposição
visível. Um filtro de vídeo exige decodificação e recodificação;
-c:v copy não pode aplicá-lo.
Para HLS, o codificador alinha os quadros-chave a cada 20 quadros
a 10 fps, correspondendo à meta de segmentos de dois segundos. event
mantém todos os segmentos, e temp_file disponibiliza cada segmento com seu
nome definitivo após o término da gravação do arquivo. A execução produz 120 quadros e então
finaliza a playlist. Trata-se de uma gravação com duração limitada de uma fonte ao vivo, não de um
acervo ilimitado com janela deslizante.
Gere um JAR executável e inicie a demonstração. --jvm system usa o JDK instalado;
--server=false evita um servidor de compilação em segundo plano. O
empacotamento --assembly do Scala CLI inclui o
ambiente de execução do Scala. -f substitui somente o
live.jar deste exemplo, e && impede a execução de
um JAR antigo quando a recompilação falha.
(
cd scala-live-demo &&
./scala-cli --power package Live.scala --assembly --preamble=false --server=false --jvm system -f -o live.jar &&
java -jar live.jar hls 60
)
O prazo limite de execução de 60 segundos começa após a validação dos argumentos e a preparação da
fonte tipográfica. Ele cobre a inicialização dos serviços, a codificação e a janela de reprodução;
a compilação e as esperas de encerramento ficam fora desse prazo. As duas portas são escolhidas
automaticamente. O terceiro e o quarto argumentos opcionais do programa selecionam as portas RTSP
e HTTP; zero solicita uma porta disponível. Se uma porta solicitada estiver ocupada, a
inicialização falha. Cada execução precisa de um novo diretório de saída diretamente dentro do
projeto; por isso, uma nova execução com hls preserva a saída existente.
Testar integrações com FFmpeg
Quando o programa imprimir a URL de reprodução, abra um segundo terminal no mesmo diretório pai e
execute este bloco. O servidor HTTP fica disponível durante a codificação e até o prazo limite da
execução. Você deverá ver um padrão em movimento em tons de cinza com uma marca de tempo branca na
faixa superior. -autoexit fecha o FFplay quando o fluxo finalizado termina;
pressione q para fechá-lo antes.
(
playback_url=$(cat scala-live-demo/hls/playback.url) &&
ffplay -autoexit "$playback_url"
)
Após a demonstração terminar com sucesso, scala-live-demo/hls contém
stream.m3u8, seis arquivos de segment-000.ts a
segment-005.ts e complete.txt. Inspecione o fluxo e decodifique a
gravação completa a partir da playlist real:
(
set -euo pipefail
cd scala-live-demo/hls &&
ffprobe -v error -show_entries stream=codec_name,width,height:format=duration \
-of json stream.m3u8 &&
ffmpeg -hide_banner -loglevel error -nostdin -xerror -progress pipe:1 -i stream.m3u8 \
-map 0:v:0 -f null - |
awk -F= '$1 == "frame" {frames=$2} END {if (frames != 120) {print "Expected 120 decoded frames" > "/dev/stderr"; exit 1}}'
)
O resultado esperado é um vídeo H.264 de 160 por 90 pixels e uma playlist com duração de 12 segundos.
A análise e a decodificação verificam a estrutura e a legibilidade da gravação; a reprodução
verifica o resultado visível do filtro. A verificação da contagem de quadros é importante porque o
FFmpeg pode pular um segmento HLS ausente e ainda retornar zero, mesmo com
-xerror. Um status de saída zero do codificador não comprova a integridade
do fluxo da câmera, a qualidade perceptual nem a latência de ponta a ponta.
Tratamento de erros e gerenciamento de recursos
O controlador verifica se o publicador e o servidor RTSP continuam ativos enquanto o codificador
trabalha. Uma desconexão pode fazer o FFmpeg terminar com menos quadros, mesmo com status de saída
zero; a verificação de progresso se recusa a marcar essa gravação como concluída. Uma falha do
codificador, uma fonte interrompida e um prazo limite atingido antes do fim da codificação retornam
status de saída um. Os logs identificam cada processo filho como mediamtx.log,
source.log ou encoder.log. Execuções de codificação que
falham podem manter mídia parcial para inspeção e não têm complete.txt.
Pressione Ctrl+C no primeiro terminal da demonstração para interrompê-la antes do prazo. O hook de encerramento da JVM bloqueia o início de novos processos filhos, encerra o HTTP e termina seus processos filhos, aguardando até dois segundos por processo antes do encerramento forçado e de uma nova espera de dois segundos. Isso se aplica a interrupções e encerramentos normais, mas não a um sinal de encerramento que não pode ser interceptado nem ao desligamento da máquina. A saída já gravada permanece; o cancelamento durante a codificação não informa a conclusão. Após a codificação ser concluída com sucesso, uma interrupção antecipada mantém a gravação completa.
Ao atingir o prazo limite normal, o programa encerra seus serviços e retorna status de saída zero se a codificação tiver sido concluída. Mantenha o diretório de saída para reproduzir novamente a playlist local ou remova esse diretório específico após a inspeção. O programa não exclui nem sobrescreve uma execução anterior, e as ferramentas baixadas e o JAR do projeto continuam disponíveis para outra execução com um novo nome de diretório.
Resolver problemas comuns
- Se a inicialização falhar, leia o log indicado e confira os pré-requisitos de executáveis e fontes. O programa encerra os processos filhos já iniciados antes de retornar um erro; ele não imprime nenhuma URL de reprodução antes que uma playlist esteja disponível.
- Se o FFplay não conseguir se conectar, inicie-o enquanto a demonstração ainda estiver em execução.
Após o encerramento, use
ffplay -autoexit scala-live-demo/hls/stream.m3u8para reproduzir novamente os arquivos locais concluídos. - Se a codificação ultrapassar o prazo limite, inspecione
encoder.loge escolha um prazo limite de até 120 segundos com um novo diretório de saída. Aumentar o prazo limite não resolve a desconexão de um publicador.
Para usar outro filtro, altere o único valor de filter, recompile o JAR e
inicie-o com um novo nome de diretório de saída.
Mantenha os limites de threads do FFmpeg enquanto mede se a codificação consegue acompanhar o
ritmo da fonte. Adicionar mais Future do Scala iniciaria mais codificadores;
isso não tornaria este codificador mais rápido nem implementaria o cancelamento desses processos.
Já para gravações locais existentes, o
guia de concatenação com Scala explica como unir arquivos
compatíveis sem recodificação.
