ZLMediaKit restful api and hook starter
115
stars
66
commits
Java
primary language
Jun 11, 2026
updated
ZLMediaKit的Spring Boot Starter,是一个针对ZLMediaKit 流媒体服务器的Java集成组件。本项目对ZLMediaKit的REST API进行了完整封装,并提供了Hook事件处理机制,支持集群化管理和多种负载均衡策略,让Java开发者能够轻松集成和管理ZLMediaKit流媒体服务器。 完整的视频平台实现:voglander
application.yml启动配置节点列表NodeSupplier接口支持动态节点发现和管理
<dependency>
<groupId>io.github.lunasaw</groupId>
<artifactId>zlm-spring-boot-starter</artifactId>
<version>${last.version}</version>
</dependency>
本项目支持两种节点配置方式:
在 application.yml 中添加ZLMediaKit配置:
zlm:
enable: true # 是否启用,未启用不会加载
balance: round_robin # 节点负载均衡算法,默认round_robin
nodes: # zlm节点列表,每个节点配置如下
- server-id: zlm-node-1 # 节点ID,可自定义
host: "http://127.0.0.1:9092" # 节点地址
secret: zlm # 节点密钥
enabled: true # 节点是否启用
hook-enabled: true # 节点是否启用hook接口,启用hook会注入hook接口,默认true,需要注意拦截器放通
- server-id: zlm-node-2 # 可配置多个节点
host: "http://127.0.0.1:9093"
secret: zlm
enabled: true
hook-enabled: true
实现NodeSupplier接口支持动态节点发现,详见动态节点发现 (NodeSupplier)章节。
import io.github.lunasaw.zlm.api.ZlmRestService;
import io.github.lunasaw.zlm.entity.ServerResponse;
import io.github.lunasaw.zlm.entity.Version;
// 获取服务器版本信息
ServerResponse<Version> versionResponse = ZlmRestService.getVersion("http://127.0.0.1:9092", "zlm");
System.out.println("ZLMediaKit版本: " + versionResponse.getData().getVersion());
// 获取流列表
ServerResponse<List<MediaData>> mediaList = ZlmRestService.getMediaList("http://127.0.0.1:9092", "zlm", new HashMap<>());
mediaList.getData().forEach(media -> {
System.out.println("流ID: " + media.getApp() + "/" + media.getStream());
});
项目内置了完整的REST API控制器,可以直接通过HTTP接口访问:
# 获取服务器版本信息
GET http://localhost:8080/zlm/api/version
# 获取流列表
POST http://localhost:8080/zlm/api/media/list
Content-Type: application/json
{
"app": "live",
"stream": ""
}
# 获取API文档
GET http://localhost:8080/swagger-ui.html
支持的API接口路径前缀:/zlm/api/,包括:
/zlm/api/version、/zlm/api/server/config/zlm/api/media/list、/zlm/api/media/close/zlm/api/proxy/add、/zlm/api/proxy/delete/zlm/api/record/start、/zlm/api/record/stop/zlm/api/rtp/open、/zlm/api/rtp/close创建Hook服务实现类来处理ZLMediaKit的事件回调:
import io.github.lunasaw.zlm.hook.service.AbstractZlmHookService;
import org.springframework.stereotype.Service;
@Service
public class CustomZlmHookService extends AbstractZlmHookService {
@Override
public HookResult onPlay(OnPlayHookParam param) {
// 播放鉴权逻辑
log.info("播放请求: {}:{}", param.getApp(), param.getStream());
if (isValidUser(param.getParams())) {
return HookResult.SUCCESS();
}
return HookResult.FAILED("无权限播放该流");
}
@Override
public HookResultForOnPublish onPublish(OnPublishHookParam param) {
// 推流鉴权逻辑
log.info("推流请求: {}:{}", param.getApp(), param.getStream());
return HookResultForOnPublish.SUCCESS();
}
@Override
public void onStreamChanged(OnStreamChangedHookParam param) {
// 流状态变化处理
log.info("流状态变化: {} - {}", param.getStream(), param.isRegist());
}
private boolean isValidUser(String params) {
// 实现用户验证逻辑
return true;
}
}
支持以下5种负载均衡算法:
| 算法 | 配置值 | 说明 |
|---|---|---|
| 随机 | random | 随机选择节点 |
| 轮询 | round_robin | 轮询选择节点(默认) |
| 一致性哈希 | consistent_hashing | 基于一致性哈希算法 |
| 加权轮询 | weight_round_robin | 基于权重的轮询 |
| 加权随机 | weight_random | 基于权重的随机选择 |
zlm:
enable: true # 是否启用ZLM功能
balance: round_robin # 负载均衡算法
nodes: # 节点配置列表
- server-id: unique-id # 节点唯一标识
host: "http://ip:port" # 节点地址
secret: "secret-key" # API密钥
enabled: true # 是否启用该节点
hook-enabled: true # 是否启用Hook功能
weight: 1 # 节点权重(仅加权算法有效)
// 获取服务器版本
ServerResponse<Version> version = ZlmRestService.getVersion(host, secret);
// 获取服务器配置
ServerResponse<ServerNodeConfig> config = ZlmRestService.getServerConfig(host, secret);
// 获取API列表
ServerResponse<List<String>> apiList = ZlmRestService.getApiList(host, secret);
// 获取服务器统计信息
ServerResponse<ImportantObjectNum> statistics = ZlmRestService.getStatistic(host, secret);
// 获取流列表
MediaReq mediaReq = new MediaReq();
mediaReq.
setApp("live");
ServerResponse<List<MediaData>> mediaList = ZlmRestService.getMediaList(host, secret, mediaReq);
// 关闭指定流
ZlmRestService.
closeStream(host, secret, mediaReq);
// 检查流是否在线
MediaOnlineStatus status = ZlmRestService.isMediaOnline(host, secret, mediaReq);
// 获取流详细信息
ServerResponse<MediaInfo> mediaInfo = ZlmRestService.getMediaInfo(host, secret, mediaReq);
// 添加拉流代理
StreamProxyItem proxyItem = new StreamProxyItem();
proxyItem.
setVhost("__defaultVhost__");
proxyItem.
setApp("live");
proxyItem.
setStream("test");
proxyItem.
setUrl("rtmp://example.com/live/stream");
ServerResponse<StreamKey> result = ZlmRestService.addStreamProxy(host, secret, proxyItem);
// 删除拉流代理
ZlmRestService.
delStreamProxy(host, secret, result.getData().
getKey());
// 添加推流
StreamPusherItem pusherItem = new StreamPusherItem();
pusherItem.
setSchema("rtmp");
pusherItem.
setVhost("__defaultVhost__");
pusherItem.
setApp("live");
pusherItem.
setStream("test");
pusherItem.
setDst_url("rtmp://push.example.com/live/stream");
ServerResponse<StreamKey> pushResult = ZlmRestService.addStreamPusherProxy(host, secret, pusherItem);
// 开始录制
RecordReq recordReq = new RecordReq();
recordReq.
setType(0); // 0-hls, 1-mp4
recordReq.
setVhost("__defaultVhost__");
recordReq.
setApp("live");
recordReq.
setStream("test");
ZlmRestService.
startRecord(host, secret, recordReq);
// 停止录制
ZlmRestService.
stopRecord(host, secret, recordReq);
// 获取录制文件
ServerResponse<Mp4RecordFile> recordFiles = ZlmRestService.getMp4RecordFile(host, secret, recordReq);
// 获取流截图
SnapshotReq snapshotReq = new SnapshotReq();
snapshotReq.
setVhost("__defaultVhost__");
snapshotReq.
setApp("live");
snapshotReq.
setStream("test");
snapshotReq.
setSavePath("/path/to/snapshot.jpg");
String result = ZlmRestService.getSnap(host, secret, snapshotReq);
// 创建RTP服务器
OpenRtpServerReq rtpReq = new OpenRtpServerReq();
rtpReq.
setStream_id("test_rtp");
rtpReq.
setPort(10000);
OpenRtpServerResult rtpResult = ZlmRestService.openRtpServer(host, secret, rtpReq);
// 关闭RTP服务器
ZlmRestService.
closeRtpServer(host, secret, "test_rtp");
实现 ZlmHookService 接口或继承 AbstractZlmHookService 类来处理各种Hook事件:
public interface ZlmHookService {
// 服务器保活事件
void onServerKeepLive(OnServerKeepaliveHookParam param);
// 播放鉴权
HookResult onPlay(OnPlayHookParam param);
// 推流鉴权
HookResultForOnPublish onPublish(OnPublishHookParam param);
// 流状态变化
void onStreamChanged(OnStreamChangedHookParam param);
// 流无人观看
HookResultForStreamNoneReader onStreamNoneReader(OnStreamNoneReaderHookParam param);
// 流未找到
void onStreamNotFound(OnStreamNotFoundHookParam param);
// 服务器启动
void onServerStarted(ServerNodeConfig param);
// RTP推流停止
void onSendRtpStopped(OnSendRtpStoppedHookParam param);
// RTP服务器超时
void onRtpServerTimeout(OnRtpServerTimeoutHookParam param);
// HTTP访问鉴权
HookResultForOnHttpAccess onHttpAccess(OnHttpAccessParam param);
// RTSP Realm鉴权
HookResultForOnRtspRealm onRtspRealm(OnRtspRealmHookParam param);
// RTSP用户密码鉴权
HookResultForOnRtspAuth onRtspAuth(OnRtspAuthHookParam param);
// 流量统计
void onFlowReport(OnFlowReportHookParam param);
// 服务器退出
void onServerExited(HookParam param);
// MP4录制完成
void onRecordMp4(OnRecordMp4HookParam param);
}
不同的Hook事件需要返回不同的结果:
// 基础Hook结果 - 用于播放鉴权等
HookResult.SUCCESS(); // 允许
HookResult.
FAILED("原因"); // 拒绝
// 推流鉴权结果
HookResultForOnPublish.
SUCCESS(); // 允许推流
HookResultForOnPublish.
FAILED("推流被拒绝"); // 拒绝推流
// HTTP访问鉴权结果
HookResultForOnHttpAccess.
SUCCESS(); // 允许访问
HookResultForOnHttpAccess.
FAILED(401,"未授权"); // 拒绝访问
本项目支持通过NodeSupplier接口实现动态节点发现和管理,支持从数据库、注册中心、配置中心等数据源动态获取节点列表。
系统默认提供DefaultNodeSupplier实现,从配置文件中获取节点列表:
@Component
public class DefaultNodeSupplier implements NodeSupplier {
@Autowired
private ZlmProperties zlmProperties;
@Override
public String getName() {
return "DefaultNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
return zlmProperties.getNodes();
}
@Override
public ZlmNode getNode(String serverId) {
return zlmProperties.getNodeMap().get(serverId);
}
}
可以实现自定义的NodeSupplier来支持动态节点发现:
@Component
public class DatabaseNodeSupplier implements NodeSupplier {
@Autowired
private NodeRepository nodeRepository;
@Override
public String getName() {
return "DatabaseNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
// 从数据库获取活跃节点列表
List<NodeEntity> activeNodes = nodeRepository.findByStatus("ACTIVE");
return activeNodes.stream()
.map(this::convertToZlmNode)
.collect(Collectors.toList());
}
@Override
public ZlmNode getNode(String serverId) {
NodeEntity entity = nodeRepository.findByServerId(serverId);
return entity != null ? convertToZlmNode(entity) : null;
}
private ZlmNode convertToZlmNode(NodeEntity entity) {
ZlmNode node = new ZlmNode();
node.setServerId(entity.getServerId());
node.setHost(entity.getHost());
node.setSecret(entity.getSecret());
node.setEnabled(entity.isEnabled());
node.setWeight(entity.getWeight());
return node;
}
}
与Spring Cloud集成,从注册中心动态发现节点:
@Component
public class EurekaNodeSupplier implements NodeSupplier {
@Autowired
private DiscoveryClient discoveryClient;
@Override
public String getName() {
return "EurekaNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
List<ServiceInstance> instances = discoveryClient.getInstances("zlm-service");
return instances.stream()
.filter(ServiceInstance::isSecure)
.map(this::convertToZlmNode)
.collect(Collectors.toList());
}
@Override
public ZlmNode getNode(String serverId) {
List<ServiceInstance> instances = discoveryClient.getInstances("zlm-service");
return instances.stream()
.filter(instance -> serverId.equals(instance.getInstanceId()))
.findFirst()
.map(this::convertToZlmNode)
.orElse(null);
}
private ZlmNode convertToZlmNode(ServiceInstance instance) {
ZlmNode node = new ZlmNode();
node.setServerId(instance.getInstanceId());
node.setHost(instance.getUri().toString());
node.setSecret(instance.getMetadata().get("secret"));
node.setEnabled(true);
node.setWeight(Integer.parseInt(instance.getMetadata().getOrDefault("weight", "1")));
return node;
}
}
从Nacos配置中心动态获取节点配置:
@Component
public class NacosNodeSupplier implements NodeSupplier {
@NacosValue("${zlm.nodes:[]}")
private String nodesConfig;
@Autowired
private ObjectMapper objectMapper;
@Override
public String getName() {
return "NacosNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
try {
if (StringUtils.hasText(nodesConfig)) {
return objectMapper.readValue(nodesConfig,
new TypeReference<List<ZlmNode>>() {
});
}
return Collections.emptyList();
} catch (Exception e) {
log.error("解析Nacos节点配置失败", e);
return Collections.emptyList();
}
}
}
DefaultNodeSupplier,配置简单zlm:
enable: true
balance: consistent_hashing
nodes:
- server-id: zlm-beijing-1
host: "http://10.0.1.10:9092"
secret: "beijing-secret"
enabled: true
hook-enabled: true
weight: 3
- server-id: zlm-beijing-2
host: "http://10.0.1.11:9092"
secret: "beijing-secret"
enabled: true
hook-enabled: true
weight: 2
- server-id: zlm-shanghai-1
host: "http://10.0.2.10:9092"
secret: "shanghai-secret"
enabled: true
hook-enabled: false # 上海节点不处理Hook
weight: 1
@Component
public class CustomLoadBalancer implements LoadBalancer {
private volatile NodeSupplier nodeSupplier;
@Override
public void setNodeSupplier(NodeSupplier nodeSupplier) {
this.nodeSupplier = nodeSupplier;
}
@Override
public ZlmNode selectNode(String key) {
List<ZlmNode> nodes = getCurrentNodes();
if (nodes == null || nodes.isEmpty()) {
return null;
}
// 实现自定义负载均衡逻辑,例如基于地理位置的选择
return selectByLocation(nodes, key);
}
@Override
public String getType() {
return "CustomLoadBalancer";
}
private List<ZlmNode> getCurrentNodes() {
if (nodeSupplier == null) {
return Collections.emptyList();
}
try {
return nodeSupplier.getNodes();
} catch (Exception e) {
log.error("获取节点列表失败", e);
return Collections.emptyList();
}
}
private ZlmNode selectByLocation(List<ZlmNode> nodes, String key) {
// 基于地理位置或其他业务逻辑的选择算法
// 例如:选择离用户最近的节点
return nodes.stream()
.filter(node -> isNearUser(node, key))
.findFirst()
.orElse(nodes.get(0));
}
private boolean isNearUser(ZlmNode node, String key) {
// 实现地理位置判断逻辑
return true;
}
}
@Service
public class ConditionalZlmHookService extends AbstractZlmHookService {
@Override
public HookResult onPlay(OnPlayHookParam param) {
// 根据不同应用进行不同处理
switch (param.getApp()) {
case "live":
return handleLivePlay(param);
case "vod":
return handleVodPlay(param);
default:
return HookResult.SUCCESS();
}
}
private HookResult handleLivePlay(OnPlayHookParam param) {
// 直播流播放逻辑
return HookResult.SUCCESS();
}
private HookResult handleVodPlay(OnPlayHookParam param) {
// 点播流播放逻辑
return HookResult.SUCCESS();
}
}
A: 请检查以下几点:
hook-enabled: trueA: 请确认:
enabled: trueA: 请检查:
@Component注解getNodes()方法返回非空且有效的节点列表A: 解决方案:
A: 建议:
A: 请检查:
getNodes()方法的性能,必要时添加缓存机制欢迎提交Issue和Pull Request!
git checkout -b feature/AmazingFeature)git commit -m 'Add some AmazingFeature')git push origin feature/AmazingFeature)本项目使用 Apache 2.0 许可证。
Java
100.0%
ZLMediaKit restful api and hook starter
115
stars
66
commits
Java
primary language
Jun 11, 2026
updated
ZLMediaKit的Spring Boot Starter,是一个针对ZLMediaKit 流媒体服务器的Java集成组件。本项目对ZLMediaKit的REST API进行了完整封装,并提供了Hook事件处理机制,支持集群化管理和多种负载均衡策略,让Java开发者能够轻松集成和管理ZLMediaKit流媒体服务器。 完整的视频平台实现:voglander
application.yml启动配置节点列表NodeSupplier接口支持动态节点发现和管理
<dependency>
<groupId>io.github.lunasaw</groupId>
<artifactId>zlm-spring-boot-starter</artifactId>
<version>${last.version}</version>
</dependency>
本项目支持两种节点配置方式:
在 application.yml 中添加ZLMediaKit配置:
zlm:
enable: true # 是否启用,未启用不会加载
balance: round_robin # 节点负载均衡算法,默认round_robin
nodes: # zlm节点列表,每个节点配置如下
- server-id: zlm-node-1 # 节点ID,可自定义
host: "http://127.0.0.1:9092" # 节点地址
secret: zlm # 节点密钥
enabled: true # 节点是否启用
hook-enabled: true # 节点是否启用hook接口,启用hook会注入hook接口,默认true,需要注意拦截器放通
- server-id: zlm-node-2 # 可配置多个节点
host: "http://127.0.0.1:9093"
secret: zlm
enabled: true
hook-enabled: true
实现NodeSupplier接口支持动态节点发现,详见动态节点发现 (NodeSupplier)章节。
import io.github.lunasaw.zlm.api.ZlmRestService;
import io.github.lunasaw.zlm.entity.ServerResponse;
import io.github.lunasaw.zlm.entity.Version;
// 获取服务器版本信息
ServerResponse<Version> versionResponse = ZlmRestService.getVersion("http://127.0.0.1:9092", "zlm");
System.out.println("ZLMediaKit版本: " + versionResponse.getData().getVersion());
// 获取流列表
ServerResponse<List<MediaData>> mediaList = ZlmRestService.getMediaList("http://127.0.0.1:9092", "zlm", new HashMap<>());
mediaList.getData().forEach(media -> {
System.out.println("流ID: " + media.getApp() + "/" + media.getStream());
});
项目内置了完整的REST API控制器,可以直接通过HTTP接口访问:
# 获取服务器版本信息
GET http://localhost:8080/zlm/api/version
# 获取流列表
POST http://localhost:8080/zlm/api/media/list
Content-Type: application/json
{
"app": "live",
"stream": ""
}
# 获取API文档
GET http://localhost:8080/swagger-ui.html
支持的API接口路径前缀:/zlm/api/,包括:
/zlm/api/version、/zlm/api/server/config/zlm/api/media/list、/zlm/api/media/close/zlm/api/proxy/add、/zlm/api/proxy/delete/zlm/api/record/start、/zlm/api/record/stop/zlm/api/rtp/open、/zlm/api/rtp/close创建Hook服务实现类来处理ZLMediaKit的事件回调:
import io.github.lunasaw.zlm.hook.service.AbstractZlmHookService;
import org.springframework.stereotype.Service;
@Service
public class CustomZlmHookService extends AbstractZlmHookService {
@Override
public HookResult onPlay(OnPlayHookParam param) {
// 播放鉴权逻辑
log.info("播放请求: {}:{}", param.getApp(), param.getStream());
if (isValidUser(param.getParams())) {
return HookResult.SUCCESS();
}
return HookResult.FAILED("无权限播放该流");
}
@Override
public HookResultForOnPublish onPublish(OnPublishHookParam param) {
// 推流鉴权逻辑
log.info("推流请求: {}:{}", param.getApp(), param.getStream());
return HookResultForOnPublish.SUCCESS();
}
@Override
public void onStreamChanged(OnStreamChangedHookParam param) {
// 流状态变化处理
log.info("流状态变化: {} - {}", param.getStream(), param.isRegist());
}
private boolean isValidUser(String params) {
// 实现用户验证逻辑
return true;
}
}
支持以下5种负载均衡算法:
| 算法 | 配置值 | 说明 |
|---|---|---|
| 随机 | random | 随机选择节点 |
| 轮询 | round_robin | 轮询选择节点(默认) |
| 一致性哈希 | consistent_hashing | 基于一致性哈希算法 |
| 加权轮询 | weight_round_robin | 基于权重的轮询 |
| 加权随机 | weight_random | 基于权重的随机选择 |
zlm:
enable: true # 是否启用ZLM功能
balance: round_robin # 负载均衡算法
nodes: # 节点配置列表
- server-id: unique-id # 节点唯一标识
host: "http://ip:port" # 节点地址
secret: "secret-key" # API密钥
enabled: true # 是否启用该节点
hook-enabled: true # 是否启用Hook功能
weight: 1 # 节点权重(仅加权算法有效)
// 获取服务器版本
ServerResponse<Version> version = ZlmRestService.getVersion(host, secret);
// 获取服务器配置
ServerResponse<ServerNodeConfig> config = ZlmRestService.getServerConfig(host, secret);
// 获取API列表
ServerResponse<List<String>> apiList = ZlmRestService.getApiList(host, secret);
// 获取服务器统计信息
ServerResponse<ImportantObjectNum> statistics = ZlmRestService.getStatistic(host, secret);
// 获取流列表
MediaReq mediaReq = new MediaReq();
mediaReq.
setApp("live");
ServerResponse<List<MediaData>> mediaList = ZlmRestService.getMediaList(host, secret, mediaReq);
// 关闭指定流
ZlmRestService.
closeStream(host, secret, mediaReq);
// 检查流是否在线
MediaOnlineStatus status = ZlmRestService.isMediaOnline(host, secret, mediaReq);
// 获取流详细信息
ServerResponse<MediaInfo> mediaInfo = ZlmRestService.getMediaInfo(host, secret, mediaReq);
// 添加拉流代理
StreamProxyItem proxyItem = new StreamProxyItem();
proxyItem.
setVhost("__defaultVhost__");
proxyItem.
setApp("live");
proxyItem.
setStream("test");
proxyItem.
setUrl("rtmp://example.com/live/stream");
ServerResponse<StreamKey> result = ZlmRestService.addStreamProxy(host, secret, proxyItem);
// 删除拉流代理
ZlmRestService.
delStreamProxy(host, secret, result.getData().
getKey());
// 添加推流
StreamPusherItem pusherItem = new StreamPusherItem();
pusherItem.
setSchema("rtmp");
pusherItem.
setVhost("__defaultVhost__");
pusherItem.
setApp("live");
pusherItem.
setStream("test");
pusherItem.
setDst_url("rtmp://push.example.com/live/stream");
ServerResponse<StreamKey> pushResult = ZlmRestService.addStreamPusherProxy(host, secret, pusherItem);
// 开始录制
RecordReq recordReq = new RecordReq();
recordReq.
setType(0); // 0-hls, 1-mp4
recordReq.
setVhost("__defaultVhost__");
recordReq.
setApp("live");
recordReq.
setStream("test");
ZlmRestService.
startRecord(host, secret, recordReq);
// 停止录制
ZlmRestService.
stopRecord(host, secret, recordReq);
// 获取录制文件
ServerResponse<Mp4RecordFile> recordFiles = ZlmRestService.getMp4RecordFile(host, secret, recordReq);
// 获取流截图
SnapshotReq snapshotReq = new SnapshotReq();
snapshotReq.
setVhost("__defaultVhost__");
snapshotReq.
setApp("live");
snapshotReq.
setStream("test");
snapshotReq.
setSavePath("/path/to/snapshot.jpg");
String result = ZlmRestService.getSnap(host, secret, snapshotReq);
// 创建RTP服务器
OpenRtpServerReq rtpReq = new OpenRtpServerReq();
rtpReq.
setStream_id("test_rtp");
rtpReq.
setPort(10000);
OpenRtpServerResult rtpResult = ZlmRestService.openRtpServer(host, secret, rtpReq);
// 关闭RTP服务器
ZlmRestService.
closeRtpServer(host, secret, "test_rtp");
实现 ZlmHookService 接口或继承 AbstractZlmHookService 类来处理各种Hook事件:
public interface ZlmHookService {
// 服务器保活事件
void onServerKeepLive(OnServerKeepaliveHookParam param);
// 播放鉴权
HookResult onPlay(OnPlayHookParam param);
// 推流鉴权
HookResultForOnPublish onPublish(OnPublishHookParam param);
// 流状态变化
void onStreamChanged(OnStreamChangedHookParam param);
// 流无人观看
HookResultForStreamNoneReader onStreamNoneReader(OnStreamNoneReaderHookParam param);
// 流未找到
void onStreamNotFound(OnStreamNotFoundHookParam param);
// 服务器启动
void onServerStarted(ServerNodeConfig param);
// RTP推流停止
void onSendRtpStopped(OnSendRtpStoppedHookParam param);
// RTP服务器超时
void onRtpServerTimeout(OnRtpServerTimeoutHookParam param);
// HTTP访问鉴权
HookResultForOnHttpAccess onHttpAccess(OnHttpAccessParam param);
// RTSP Realm鉴权
HookResultForOnRtspRealm onRtspRealm(OnRtspRealmHookParam param);
// RTSP用户密码鉴权
HookResultForOnRtspAuth onRtspAuth(OnRtspAuthHookParam param);
// 流量统计
void onFlowReport(OnFlowReportHookParam param);
// 服务器退出
void onServerExited(HookParam param);
// MP4录制完成
void onRecordMp4(OnRecordMp4HookParam param);
}
不同的Hook事件需要返回不同的结果:
// 基础Hook结果 - 用于播放鉴权等
HookResult.SUCCESS(); // 允许
HookResult.
FAILED("原因"); // 拒绝
// 推流鉴权结果
HookResultForOnPublish.
SUCCESS(); // 允许推流
HookResultForOnPublish.
FAILED("推流被拒绝"); // 拒绝推流
// HTTP访问鉴权结果
HookResultForOnHttpAccess.
SUCCESS(); // 允许访问
HookResultForOnHttpAccess.
FAILED(401,"未授权"); // 拒绝访问
本项目支持通过NodeSupplier接口实现动态节点发现和管理,支持从数据库、注册中心、配置中心等数据源动态获取节点列表。
系统默认提供DefaultNodeSupplier实现,从配置文件中获取节点列表:
@Component
public class DefaultNodeSupplier implements NodeSupplier {
@Autowired
private ZlmProperties zlmProperties;
@Override
public String getName() {
return "DefaultNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
return zlmProperties.getNodes();
}
@Override
public ZlmNode getNode(String serverId) {
return zlmProperties.getNodeMap().get(serverId);
}
}
可以实现自定义的NodeSupplier来支持动态节点发现:
@Component
public class DatabaseNodeSupplier implements NodeSupplier {
@Autowired
private NodeRepository nodeRepository;
@Override
public String getName() {
return "DatabaseNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
// 从数据库获取活跃节点列表
List<NodeEntity> activeNodes = nodeRepository.findByStatus("ACTIVE");
return activeNodes.stream()
.map(this::convertToZlmNode)
.collect(Collectors.toList());
}
@Override
public ZlmNode getNode(String serverId) {
NodeEntity entity = nodeRepository.findByServerId(serverId);
return entity != null ? convertToZlmNode(entity) : null;
}
private ZlmNode convertToZlmNode(NodeEntity entity) {
ZlmNode node = new ZlmNode();
node.setServerId(entity.getServerId());
node.setHost(entity.getHost());
node.setSecret(entity.getSecret());
node.setEnabled(entity.isEnabled());
node.setWeight(entity.getWeight());
return node;
}
}
与Spring Cloud集成,从注册中心动态发现节点:
@Component
public class EurekaNodeSupplier implements NodeSupplier {
@Autowired
private DiscoveryClient discoveryClient;
@Override
public String getName() {
return "EurekaNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
List<ServiceInstance> instances = discoveryClient.getInstances("zlm-service");
return instances.stream()
.filter(ServiceInstance::isSecure)
.map(this::convertToZlmNode)
.collect(Collectors.toList());
}
@Override
public ZlmNode getNode(String serverId) {
List<ServiceInstance> instances = discoveryClient.getInstances("zlm-service");
return instances.stream()
.filter(instance -> serverId.equals(instance.getInstanceId()))
.findFirst()
.map(this::convertToZlmNode)
.orElse(null);
}
private ZlmNode convertToZlmNode(ServiceInstance instance) {
ZlmNode node = new ZlmNode();
node.setServerId(instance.getInstanceId());
node.setHost(instance.getUri().toString());
node.setSecret(instance.getMetadata().get("secret"));
node.setEnabled(true);
node.setWeight(Integer.parseInt(instance.getMetadata().getOrDefault("weight", "1")));
return node;
}
}
从Nacos配置中心动态获取节点配置:
@Component
public class NacosNodeSupplier implements NodeSupplier {
@NacosValue("${zlm.nodes:[]}")
private String nodesConfig;
@Autowired
private ObjectMapper objectMapper;
@Override
public String getName() {
return "NacosNodeSupplier";
}
@Override
public List<ZlmNode> getNodes() {
try {
if (StringUtils.hasText(nodesConfig)) {
return objectMapper.readValue(nodesConfig,
new TypeReference<List<ZlmNode>>() {
});
}
return Collections.emptyList();
} catch (Exception e) {
log.error("解析Nacos节点配置失败", e);
return Collections.emptyList();
}
}
}
DefaultNodeSupplier,配置简单zlm:
enable: true
balance: consistent_hashing
nodes:
- server-id: zlm-beijing-1
host: "http://10.0.1.10:9092"
secret: "beijing-secret"
enabled: true
hook-enabled: true
weight: 3
- server-id: zlm-beijing-2
host: "http://10.0.1.11:9092"
secret: "beijing-secret"
enabled: true
hook-enabled: true
weight: 2
- server-id: zlm-shanghai-1
host: "http://10.0.2.10:9092"
secret: "shanghai-secret"
enabled: true
hook-enabled: false # 上海节点不处理Hook
weight: 1
@Component
public class CustomLoadBalancer implements LoadBalancer {
private volatile NodeSupplier nodeSupplier;
@Override
public void setNodeSupplier(NodeSupplier nodeSupplier) {
this.nodeSupplier = nodeSupplier;
}
@Override
public ZlmNode selectNode(String key) {
List<ZlmNode> nodes = getCurrentNodes();
if (nodes == null || nodes.isEmpty()) {
return null;
}
// 实现自定义负载均衡逻辑,例如基于地理位置的选择
return selectByLocation(nodes, key);
}
@Override
public String getType() {
return "CustomLoadBalancer";
}
private List<ZlmNode> getCurrentNodes() {
if (nodeSupplier == null) {
return Collections.emptyList();
}
try {
return nodeSupplier.getNodes();
} catch (Exception e) {
log.error("获取节点列表失败", e);
return Collections.emptyList();
}
}
private ZlmNode selectByLocation(List<ZlmNode> nodes, String key) {
// 基于地理位置或其他业务逻辑的选择算法
// 例如:选择离用户最近的节点
return nodes.stream()
.filter(node -> isNearUser(node, key))
.findFirst()
.orElse(nodes.get(0));
}
private boolean isNearUser(ZlmNode node, String key) {
// 实现地理位置判断逻辑
return true;
}
}
@Service
public class ConditionalZlmHookService extends AbstractZlmHookService {
@Override
public HookResult onPlay(OnPlayHookParam param) {
// 根据不同应用进行不同处理
switch (param.getApp()) {
case "live":
return handleLivePlay(param);
case "vod":
return handleVodPlay(param);
default:
return HookResult.SUCCESS();
}
}
private HookResult handleLivePlay(OnPlayHookParam param) {
// 直播流播放逻辑
return HookResult.SUCCESS();
}
private HookResult handleVodPlay(OnPlayHookParam param) {
// 点播流播放逻辑
return HookResult.SUCCESS();
}
}
A: 请检查以下几点:
hook-enabled: trueA: 请确认:
enabled: trueA: 请检查:
@Component注解getNodes()方法返回非空且有效的节点列表A: 解决方案:
A: 建议:
A: 请检查:
getNodes()方法的性能,必要时添加缓存机制欢迎提交Issue和Pull Request!
git checkout -b feature/AmazingFeature)git commit -m 'Add some AmazingFeature')git push origin feature/AmazingFeature)本项目使用 Apache 2.0 许可证。
Java
100.0%