FFmpeg avec Scala : vidéo en temps réel
Utilisez Scala pour démarrer un flux RTSP de test en direct, le transformer avec FFmpeg et servir le flux HLS obtenu via HTTP en local. Vous verrez un motif animé en niveaux de gris avec un horodatage et conserverez un enregistrement de 12 secondes. L’exemple fonctionne sous Linux, ne nécessite aucune caméra et arrête ses propres services à la fin de la démonstration.
L’enjeu essentiel est la maîtrise des processus lancés : Scala lance FFmpeg et attend sa fin
effective, tout en conservant un délai limite et un hook d’arrêt capables de l’arrêter. Attendre
uniquement un Future Scala n’annule pas un encodeur externe. Il s’agit d’un
seul flux vidéo, sans audio, avec un seul niveau de qualité et une mise en tampon HLS classique ;
cet exemple ne démontre ni objectif de latence en production ni diffusion à débit adaptatif.
Configurer Scala et FFmpeg
Ce tutoriel a été testé sous Linux x86-64 avec OpenJDK 21.0.12.1, Scala CLI 1.17.1,
Scala 2.13.18, MediaMTX 1.21.1 et
FFmpeg/ffprobe/ffplay 9.0.1. Utilisez Bash pour les commandes. Installez une
version de FFmpeg pour Linux avec libx264,
le multiplexeur HLS et les filtres testsrc2, scale,
hue, drawbox et drawtext.
Vous avez également besoin de ce JDK, de cURL, de gzip, de tar, d’awk, de
fc-match de fontconfig et d’une police TrueType installée, comme Liberation Sans.
FFplay nécessite un environnement de bureau graphique pour afficher la vidéo.
java -version && ffmpeg -version && ffprobe -version && ffplay -version &&
fc-match 'Liberation Sans'
Créez un nouveau projet et téléchargez les versions fixées de
Scala CLI et de
MediaMTX. MediaMTX est le point de terminaison RTSP local :
FFmpeg y publie le flux, puis un autre processus FFmpeg lit ce flux.
Collez ce bloc depuis un répertoire parent accessible en écriture. Si
scala-live-demo existe déjà, la configuration s’arrête sans modifier ce répertoire ;
choisissez un autre répertoire parent pour recommencer.
(
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
)
Restez dans le répertoire parent pour les commandes suivantes. Chaque bloc entre parenthèses limite
ses changements de répertoire et ses options de shell à son propre contexte. Un échec de
téléchargement peut laisser un nouveau projet incomplet ; inspectez-le et utilisez un autre
répertoire parent avant de réessayer. Enregistrez le programme complet ci-dessous dans
scala-live-demo/Live.scala. La première compilation télécharge le compilateur Scala et les
bibliothèques ; il n’y a aucune dépendance applicative supplémentaire. Les commandes ciblent
explicitement Live.scala ; un éventuel projet sbt englobant n’est pas utilisé
comme entrée de compilation.
Intégrer FFmpeg à Scala
Le programme utilise ProcessBuilder de Java avec des arguments séparés, plutôt
que d’interpoler une commande shell. Il écrit une configuration MediaMTX privée, associe le serveur
HTTP à l’interface de boucle locale et choisit des ports disponibles par défaut. Seuls la playlist
et les fichiers de segments terminés sont servis ; la configuration, les polices et les journaux
ne sont pas exposés via 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) }
}
}
Traitement de flux vidéo en temps réel
L’option -re de FFmpeg cadence la source générée au lieu de la publier
aussi vite que le processeur peut l’encoder. MediaMTX accepte le processus de publication sur
/live ; l’encodeur lit ce même
chemin RTSP en utilisant le
transport TCP. Cet exemple utilise RTSP, sans supposer que RTMP
ou une caméra quelconque présente le même comportement en entrée.
Le filtre redimensionne la vidéo à 160 par 90 pixels, supprime les couleurs et dessine une bande
noire avec un horodatage du temps écoulé.
drawtext nécessite la prise en charge des
polices dans votre version de FFmpeg, même si le fichier de police est indiqué explicitement.
Les versions intégrant fontconfig peuvent utiliser une autre police si ce fichier est invalide ;
utilisez donc la police installée sélectionnée lors de la configuration et vérifiez l’incrustation
visible. Un filtre vidéo nécessite un décodage et un réencodage ;
-c:v copy ne peut pas l’appliquer.
Pour le HLS, l’encodeur aligne les images clés sur un intervalle de
20 images à 10 fps, ce qui correspond à la durée cible de deux secondes par segment.
event conserve tous les segments, et temp_file
rend chaque segment disponible sous son nom définitif une fois l’écriture terminée. L’exécution
produit 120 images, puis finalise la playlist. Il s’agit d’un enregistrement de durée limitée
d’une source en direct, et non d’une archive illimitée à fenêtre glissante.
Compilez un JAR exécutable et démarrez la démonstration. --jvm system utilise
votre JDK installé ; --server=false évite un serveur de compilation en arrière-plan.
La création de paquet avec --assembly de Scala CLI
inclut l’environnement d’exécution Scala. -f ne remplace que
live.jar de cet exemple, et && empêche l’exécution
d’un ancien JAR en cas d’échec de la recompilation.
(
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
)
Le délai limite d’exécution de 60 secondes commence après la validation des arguments et la mise en
place de la police. Il couvre le démarrage des services, l’encodage et la fenêtre de lecture ;
la compilation et les attentes lors de l’arrêt n’y sont pas incluses. Les deux ports sont choisis
automatiquement. Les troisième et quatrième arguments du programme, facultatifs, permettent de
sélectionner les ports RTSP et HTTP ; zéro demande un port disponible. Si un port demandé est déjà
occupé, le démarrage échoue. Chaque exécution nécessite un nouveau répertoire de sortie directement
dans le projet ; une nouvelle exécution avec hls préserve donc la sortie
existante.
Tester les intégrations de FFmpeg
Lorsque le programme affiche l’URL de lecture, ouvrez un second terminal dans le même répertoire
parent et exécutez ce bloc. Le serveur HTTP est disponible pendant l’encodage et jusqu’à
l’expiration du délai limite d’exécution. Vous devriez voir un motif animé en niveaux de gris avec
un horodatage blanc dans la bande supérieure. -autoexit ferme FFplay lorsque
le flux finalisé se termine ; appuyez sur q pour le fermer plus tôt.
(
playback_url=$(cat scala-live-demo/hls/playback.url) &&
ffplay -autoexit "$playback_url"
)
Une fois la démonstration terminée avec succès, scala-live-demo/hls contient
stream.m3u8, les six fichiers de segment-000.ts à
segment-005.ts, ainsi que complete.txt. Inspectez le flux et
décodez l’enregistrement complet à partir de sa playlist réelle :
(
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}}'
)
Vous devriez obtenir une vidéo H.264 de 160 par 90 pixels et une playlist d’une durée de
12 secondes. L’analyse et le décodage vérifient la structure et la lisibilité de l’enregistrement ;
la lecture vérifie le résultat visible du filtre. La vérification du nombre d’images est
importante, car FFmpeg peut ignorer un segment HLS manquant tout en renvoyant zéro, même avec
-xerror. Un code de retour nul de l’encodeur ne permet de conclure ni à
l’intégrité du flux d’une caméra, ni à sa qualité perceptuelle, ni à sa latence de bout en bout.
Gestion des erreurs et des ressources
Le contrôleur vérifie que le processus de publication et le serveur RTSP restent actifs pendant
que l’encodeur travaille. Une déconnexion peut amener FFmpeg à terminer avec moins d’images, même
avec un code de retour nul ; la vérification de la progression refuse alors de marquer cet
enregistrement comme terminé. Un échec de l’encodeur, une source arrêtée et l’expiration du délai
limite avant la fin de l’encodage renvoient chacun un code de retour égal à un. Les journaux
identifient chaque processus enfant par mediamtx.log,
source.log ou encoder.log. Les exécutions dont l’encodage
échoue peuvent conserver des médias partiels à des fins d’inspection et ne contiennent pas de
complete.txt.
Appuyez sur Ctrl+C dans le premier terminal de la démonstration pour l’arrêter prématurément. Le hook d’arrêt de la JVM empêche le lancement de nouveaux processus enfants, arrête HTTP et met fin à ses processus enfants, en attendant jusqu’à deux secondes par processus avant un arrêt forcé, puis en attendant encore deux secondes. Ce mécanisme s’applique aux interruptions et aux arrêts ordinaires, mais pas à un signal d’arrêt non interceptable ni à l’arrêt de la machine. Les données de sortie déjà écrites sont conservées ; en cas d’annulation pendant l’encodage, l’enregistrement n’est pas signalé comme terminé. Une fois l’encodage réussi, un arrêt anticipé conserve l’enregistrement terminé.
À l’expiration normale du délai, le programme arrête ses services et renvoie un code de retour nul si l’encodage est terminé. Conservez le répertoire de sortie pour relire la playlist locale, ou supprimez ce répertoire précis après inspection. Le programme ne supprime ni n’écrase les résultats d’une exécution antérieure, et les outils téléchargés ainsi que le JAR du projet restent disponibles pour une nouvelle exécution avec un nouveau nom de répertoire.
Résoudre les problèmes courants
- Si le démarrage échoue, lisez le journal indiqué et vérifiez les prérequis concernant l’exécutable et la police. Le programme arrête les processus enfants déjà lancés avant de renvoyer une erreur ; il n’affiche aucune URL de lecture tant qu’aucune playlist n’est disponible.
- Si FFplay ne peut pas se connecter, démarrez-le pendant que la démonstration est encore en cours.
Après l’arrêt, utilisez
ffplay -autoexit scala-live-demo/hls/stream.m3u8pour relire les fichiers locaux terminés. - Si l’encodage dépasse le délai limite, consultez
encoder.loget choisissez un délai allant jusqu’à 120 secondes avec un nouveau répertoire de sortie. Allonger le délai ne corrige pas la déconnexion du processus de publication.
Pour utiliser un autre filtre, modifiez l’unique valeur filter, recompilez
le JAR et lancez-le avec un nouveau nom de répertoire de sortie.
Conservez les limites de threads de FFmpeg lorsque vous mesurez la capacité de l’encodage à suivre
le rythme de la source. Ajouter des Future Scala lancerait davantage
d’encodeurs ; cela ne rendrait pas cet encodeur plus rapide et ne permettrait pas d’annuler ces
processus. Pour les enregistrements locaux existants, le
guide de concaténation avec Scala explique plutôt comment assembler
des fichiers compatibles sans réencodage.
