mirror of
https://gitee.com/ulthon/ulthon_admin.git
synced 2026-08-30 12:45:32 +08:00
feat(nginx-log): 多节点适配(node_id 全链路 + 全局/按节点聚合 + 节点筛选 + v2.4.0 升级脚本)
- 6 张表加 node_id 字段 + position/stat 改复合唯一键(migration) - 6 个 Scheme 同步 node_id 注解 + 唯一键 - Reader 带 node_id 参数 + 懒加载接管(getLastPosition 回退查 node_id='') - Aggregator 全局行(node_id='')+ 按节点循环 insertAggregatesForScope - Stat 所有查询方法加 where node_id 条件 - Dashboard 控制器/视图加节点筛选下拉框 + AJAX 带 node_id - AccessLog 列表加节点列 + 采集进度卡片显示节点信息 - import 定时任务 run_type 改 all(每节点采集自己的日志) - v2.4.0 升级脚本追加 ALTER TABLE(幂等 check-then-alter) - 菜单调整:去掉 Nginx 顶级菜单 + 读取位置管理,改为系统管理下挂两个子菜单 - ulthon-timer 文档补充 run_type=all 多节点部署注意事项
This commit is contained in:
@@ -24,9 +24,9 @@ class NginxLogAggregatorServiceBase
|
||||
/**
|
||||
* 聚合指定日期的指定小时.
|
||||
*
|
||||
* 流程:
|
||||
* 1. stat_hour:DELETE WHERE stat_date=:d AND stat_hour=:h → INSERT 1 行(该小时聚合)
|
||||
* 2. stat_url / stat_referer / stat_ua:DELETE WHERE stat_date=:d(整天重跑)→ INSERT top N 行
|
||||
* 流程(每张表生成 全局行 node_id='' + 各节点行):
|
||||
* 1. stat_hour:DELETE WHERE stat_date=:d AND stat_hour=:h → INSERT 全局 1 行 + 每节点 1 行
|
||||
* 2. stat_url / stat_referer / stat_ua:DELETE WHERE stat_date=:d(整天重跑)→ INSERT 全局 topN + 每节点 topN
|
||||
*
|
||||
* @param int $statDate 日期 YYYYMMDD(如 20260728)
|
||||
* @param int $statHour 小时 0-23
|
||||
@@ -65,142 +65,50 @@ class NginxLogAggregatorServiceBase
|
||||
|
||||
Db::startTrans();
|
||||
try {
|
||||
// 1. stat_hour:DELETE + INSERT 1 行
|
||||
// DELETE(逻辑不变):stat_hour 按小时,stat_url/referer/ua 按整天
|
||||
NginxStatHour::where('stat_date', $statDate)
|
||||
->where('stat_hour', $statHour)
|
||||
->delete();
|
||||
|
||||
$hourSql = "INSERT INTO `{$hourTable}`
|
||||
(stat_date, stat_hour, pv, uv, total_bytes, status_2xx, status_3xx, status_4xx, status_5xx, status_other, avg_request_time)
|
||||
SELECT :stat_date, :stat_hour,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv,
|
||||
COALESCE(SUM(body_bytes_sent), 0) AS total_bytes,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 200 AND 299 THEN 1 ELSE 0 END), 0) AS status_2xx,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 300 AND 399 THEN 1 ELSE 0 END), 0) AS status_3xx,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 400 AND 499 THEN 1 ELSE 0 END), 0) AS status_4xx,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 500 AND 599 THEN 1 ELSE 0 END), 0) AS status_5xx,
|
||||
COALESCE(SUM(CASE WHEN status < 200 OR status > 599 THEN 1 ELSE 0 END), 0) AS status_other,
|
||||
COALESCE(AVG(request_time), 0) AS avg_request_time
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :hour_start AND time_local < :hour_end
|
||||
HAVING COUNT(*) > 0";
|
||||
|
||||
$statHourRows = (int) Db::execute($hourSql, [
|
||||
'stat_date' => $statDate,
|
||||
'stat_hour' => $statHour,
|
||||
'hour_start' => $hourStart,
|
||||
'hour_end' => $hourEnd,
|
||||
]);
|
||||
|
||||
// 2. stat_url:DELETE(整天)+ INSERT topN
|
||||
NginxStatUrl::where('stat_date', $statDate)->delete();
|
||||
|
||||
$urlLimit = $topn > 0 ? 'LIMIT ' . $topn : '';
|
||||
$urlSql = "INSERT INTO `{$urlTable}`
|
||||
(stat_date, uri, pv, uv, total_bytes, avg_request_time)
|
||||
SELECT :stat_date, uri,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv,
|
||||
COALESCE(SUM(body_bytes_sent), 0) AS total_bytes,
|
||||
COALESCE(AVG(request_time), 0) AS avg_request_time
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :day_start AND time_local < :day_end
|
||||
GROUP BY uri
|
||||
ORDER BY pv DESC
|
||||
{$urlLimit}";
|
||||
|
||||
$statUrlRows = (int) Db::execute($urlSql, [
|
||||
'stat_date' => $statDate,
|
||||
'day_start' => $dayStart,
|
||||
'day_end' => $dayEnd,
|
||||
]);
|
||||
|
||||
// 3. stat_referer:DELETE + INSERT topN(从 http_referer 提取 domain)
|
||||
NginxStatReferer::where('stat_date', $statDate)->delete();
|
||||
|
||||
$refererLimit = $topn > 0 ? 'LIMIT ' . $topn : '';
|
||||
$refererSql = "INSERT INTO `{$refererTable}`
|
||||
(stat_date, referer_domain, pv, uv)
|
||||
SELECT :stat_date,
|
||||
CASE
|
||||
WHEN http_referer IS NULL OR http_referer = '' OR http_referer = '-' THEN '-'
|
||||
ELSE SUBSTRING_INDEX(
|
||||
SUBSTRING_INDEX(
|
||||
REPLACE(REPLACE(http_referer, 'https://', ''), 'http://', ''),
|
||||
'/', 1
|
||||
),
|
||||
'?', 1
|
||||
)
|
||||
END AS referer_domain,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :day_start AND time_local < :day_end
|
||||
GROUP BY referer_domain
|
||||
ORDER BY pv DESC
|
||||
{$refererLimit}";
|
||||
|
||||
$statRefererRows = (int) Db::execute($refererSql, [
|
||||
'stat_date' => $statDate,
|
||||
'day_start' => $dayStart,
|
||||
'day_end' => $dayEnd,
|
||||
]);
|
||||
|
||||
// 4. stat_ua:DELETE + INSERT topN(从 http_user_agent 分类)
|
||||
NginxStatUa::where('stat_date', $statDate)->delete();
|
||||
|
||||
$uaLimit = $topn > 0 ? 'LIMIT ' . $topn : '';
|
||||
$uaSql = "INSERT INTO `{$uaTable}`
|
||||
(stat_date, ua_type, ua_name, pv, uv)
|
||||
SELECT :stat_date,
|
||||
CASE
|
||||
WHEN http_user_agent LIKE '%bot%'
|
||||
OR http_user_agent LIKE '%spider%'
|
||||
OR http_user_agent LIKE '%crawl%'
|
||||
OR http_user_agent LIKE '%slurp%'
|
||||
OR http_user_agent LIKE '%bingpreview%'
|
||||
OR http_user_agent LIKE '%facebookexternalhit%'
|
||||
OR http_user_agent LIKE '%twitterbot%' THEN 'spider'
|
||||
WHEN http_user_agent IS NULL OR http_user_agent = '' OR http_user_agent = '-' THEN 'unknown'
|
||||
WHEN http_user_agent LIKE '%Mozilla%'
|
||||
OR http_user_agent LIKE '%Chrome%'
|
||||
OR http_user_agent LIKE '%Safari%'
|
||||
OR http_user_agent LIKE '%Firefox%'
|
||||
OR http_user_agent LIKE '%Edg%'
|
||||
OR http_user_agent LIKE '%Opera%'
|
||||
OR http_user_agent LIKE '%MSIE%'
|
||||
OR http_user_agent LIKE '%Trident%' THEN 'browser'
|
||||
ELSE 'unknown'
|
||||
END AS ua_type,
|
||||
CASE
|
||||
WHEN http_user_agent LIKE '%Googlebot%' THEN 'Googlebot'
|
||||
WHEN http_user_agent LIKE '%Baiduspider%' THEN 'Baiduspider'
|
||||
WHEN http_user_agent LIKE '%bingbot%' THEN 'Bingbot'
|
||||
WHEN http_user_agent LIKE '%DuckDuckBot%' THEN 'DuckDuckBot'
|
||||
WHEN http_user_agent LIKE '%YandexBot%' THEN 'YandexBot'
|
||||
WHEN http_user_agent LIKE '%Edg/%' THEN 'Microsoft Edge'
|
||||
WHEN http_user_agent LIKE '%OPR/%' OR http_user_agent LIKE '%Opera%' THEN 'Opera'
|
||||
WHEN http_user_agent LIKE '%Firefox/%' THEN 'Firefox'
|
||||
WHEN http_user_agent LIKE '%Chrome/%' THEN 'Chrome'
|
||||
WHEN http_user_agent LIKE '%Safari/%' THEN 'Safari'
|
||||
WHEN http_user_agent LIKE '%MSIE%' OR http_user_agent LIKE '%Trident%' THEN 'Internet Explorer'
|
||||
WHEN http_user_agent IS NULL OR http_user_agent = '' OR http_user_agent = '-' THEN 'Unknown'
|
||||
ELSE 'Other'
|
||||
END AS ua_name,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :day_start AND time_local < :day_end
|
||||
GROUP BY ua_type, ua_name
|
||||
ORDER BY pv DESC
|
||||
{$uaLimit}";
|
||||
$statHourRows = 0;
|
||||
$statUrlRows = 0;
|
||||
$statRefererRows = 0;
|
||||
$statUaRows = 0;
|
||||
|
||||
$statUaRows = (int) Db::execute($uaSql, [
|
||||
'stat_date' => $statDate,
|
||||
'day_start' => $dayStart,
|
||||
'day_end' => $dayEnd,
|
||||
]);
|
||||
// 批次 A:全局行(node_id='',跨所有节点聚合,UV 精确)
|
||||
$global = $this->insertAggregatesForScope(
|
||||
'', null,
|
||||
$statDate, $statHour, $hourStart, $hourEnd, $dayStart, $dayEnd, $topn,
|
||||
$rawTable, $hourTable, $urlTable, $refererTable, $uaTable
|
||||
);
|
||||
$statHourRows += $global['stat_hour'];
|
||||
$statUrlRows += $global['stat_url'];
|
||||
$statRefererRows += $global['stat_referer'];
|
||||
$statUaRows += $global['stat_ua'];
|
||||
|
||||
// 批次 B:按节点行(循环每个 node_id)
|
||||
// 用 day 窗口取节点列表:保证 day 窗口表(url/referer/ua)的 per-node 行不会因
|
||||
// 某节点"本小时无活动"而被 DELETE 后丢失;stat_hour 的 per-node 行由
|
||||
// HAVING COUNT(*)>0 自然过滤掉本小时无活动的节点。
|
||||
$nodeRows = Db::query(
|
||||
"SELECT DISTINCT node_id FROM `{$rawTable}` WHERE time_local >= :day_start AND time_local < :day_end AND node_id != ''",
|
||||
['day_start' => $dayStart, 'day_end' => $dayEnd]
|
||||
);
|
||||
foreach ($nodeRows as $nodeRow) {
|
||||
$nid = (string) $nodeRow['node_id'];
|
||||
$perNode = $this->insertAggregatesForScope(
|
||||
$nid, $nid,
|
||||
$statDate, $statHour, $hourStart, $hourEnd, $dayStart, $dayEnd, $topn,
|
||||
$rawTable, $hourTable, $urlTable, $refererTable, $uaTable
|
||||
);
|
||||
$statHourRows += $perNode['stat_hour'];
|
||||
$statUrlRows += $perNode['stat_url'];
|
||||
$statRefererRows += $perNode['stat_referer'];
|
||||
$statUaRows += $perNode['stat_ua'];
|
||||
}
|
||||
|
||||
Db::commit();
|
||||
|
||||
@@ -216,6 +124,175 @@ class NginxLogAggregatorServiceBase
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 为单个作用域(全局或某节点)执行 4 段 INSERT.
|
||||
*
|
||||
* - 全局作用域:$nodeId='' 且 $nodeFilter=null(不追加 node WHERE,跨所有节点聚合)
|
||||
* - 节点作用域:$nodeId=节点ID 且 $nodeFilter=节点ID(追加 AND node_id=:node_filter)
|
||||
*
|
||||
* stat_hour:1 行(无 GROUP BY,HAVING COUNT(*)>0)
|
||||
* stat_url / stat_referer / stat_ua:topN(各自 GROUP BY + LIMIT)
|
||||
*
|
||||
* @param string $nodeId INSERT 的 node_id 列值(全局为 '')
|
||||
* @param string|null $nodeFilter null=全局(不追加 WHERE);非 null=追加节点过滤
|
||||
* @return array {stat_hour:int, stat_url:int, stat_referer:int, stat_ua:int}
|
||||
*/
|
||||
protected function insertAggregatesForScope(
|
||||
string $nodeId,
|
||||
?string $nodeFilter,
|
||||
int $statDate,
|
||||
int $statHour,
|
||||
int $hourStart,
|
||||
int $hourEnd,
|
||||
int $dayStart,
|
||||
int $dayEnd,
|
||||
int $topn,
|
||||
string $rawTable,
|
||||
string $hourTable,
|
||||
string $urlTable,
|
||||
string $refererTable,
|
||||
string $uaTable
|
||||
): array {
|
||||
// 节点过滤片段:全局行无;节点行追加 AND node_id = :node_filter
|
||||
$nodeWhere = $nodeFilter === null ? '' : 'AND node_id = :node_filter';
|
||||
|
||||
// hour 窗口绑定(stat_hour 专用)
|
||||
$hourParams = [
|
||||
'node_id' => $nodeId,
|
||||
'stat_date' => $statDate,
|
||||
'stat_hour' => $statHour,
|
||||
'hour_start' => $hourStart,
|
||||
'hour_end' => $hourEnd,
|
||||
];
|
||||
if ($nodeFilter !== null) {
|
||||
$hourParams['node_filter'] = $nodeFilter;
|
||||
}
|
||||
|
||||
// day 窗口绑定(stat_url/referer/ua 共用)
|
||||
$dayParams = [
|
||||
'node_id' => $nodeId,
|
||||
'stat_date' => $statDate,
|
||||
'day_start' => $dayStart,
|
||||
'day_end' => $dayEnd,
|
||||
];
|
||||
if ($nodeFilter !== null) {
|
||||
$dayParams['node_filter'] = $nodeFilter;
|
||||
}
|
||||
|
||||
// 1. stat_hour:1 行(无 GROUP BY)
|
||||
$hourSql = "INSERT INTO `{$hourTable}`
|
||||
(node_id, stat_date, stat_hour, pv, uv, total_bytes, status_2xx, status_3xx, status_4xx, status_5xx, status_other, avg_request_time)
|
||||
SELECT :node_id, :stat_date, :stat_hour,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv,
|
||||
COALESCE(SUM(body_bytes_sent), 0) AS total_bytes,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 200 AND 299 THEN 1 ELSE 0 END), 0) AS status_2xx,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 300 AND 399 THEN 1 ELSE 0 END), 0) AS status_3xx,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 400 AND 499 THEN 1 ELSE 0 END), 0) AS status_4xx,
|
||||
COALESCE(SUM(CASE WHEN status BETWEEN 500 AND 599 THEN 1 ELSE 0 END), 0) AS status_5xx,
|
||||
COALESCE(SUM(CASE WHEN status < 200 OR status > 599 THEN 1 ELSE 0 END), 0) AS status_other,
|
||||
COALESCE(AVG(request_time), 0) AS avg_request_time
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :hour_start AND time_local < :hour_end {$nodeWhere}
|
||||
HAVING COUNT(*) > 0";
|
||||
$statHourRows = (int) Db::execute($hourSql, $hourParams);
|
||||
|
||||
// 2. stat_url:topN(GROUP BY uri)
|
||||
$urlLimit = $topn > 0 ? 'LIMIT ' . $topn : '';
|
||||
$urlSql = "INSERT INTO `{$urlTable}`
|
||||
(node_id, stat_date, uri, pv, uv, total_bytes, avg_request_time)
|
||||
SELECT :node_id, :stat_date, uri,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv,
|
||||
COALESCE(SUM(body_bytes_sent), 0) AS total_bytes,
|
||||
COALESCE(AVG(request_time), 0) AS avg_request_time
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :day_start AND time_local < :day_end {$nodeWhere}
|
||||
GROUP BY uri
|
||||
ORDER BY pv DESC
|
||||
{$urlLimit}";
|
||||
$statUrlRows = (int) Db::execute($urlSql, $dayParams);
|
||||
|
||||
// 3. stat_referer:topN(从 http_referer 提取 domain)
|
||||
$refererLimit = $topn > 0 ? 'LIMIT ' . $topn : '';
|
||||
$refererSql = "INSERT INTO `{$refererTable}`
|
||||
(node_id, stat_date, referer_domain, pv, uv)
|
||||
SELECT :node_id, :stat_date,
|
||||
CASE
|
||||
WHEN http_referer IS NULL OR http_referer = '' OR http_referer = '-' THEN '-'
|
||||
ELSE SUBSTRING_INDEX(
|
||||
SUBSTRING_INDEX(
|
||||
REPLACE(REPLACE(http_referer, 'https://', ''), 'http://', ''),
|
||||
'/', 1
|
||||
),
|
||||
'?', 1
|
||||
)
|
||||
END AS referer_domain,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :day_start AND time_local < :day_end {$nodeWhere}
|
||||
GROUP BY referer_domain
|
||||
ORDER BY pv DESC
|
||||
{$refererLimit}";
|
||||
$statRefererRows = (int) Db::execute($refererSql, $dayParams);
|
||||
|
||||
// 4. stat_ua:topN(从 http_user_agent 分类)
|
||||
$uaLimit = $topn > 0 ? 'LIMIT ' . $topn : '';
|
||||
$uaSql = "INSERT INTO `{$uaTable}`
|
||||
(node_id, stat_date, ua_type, ua_name, pv, uv)
|
||||
SELECT :node_id, :stat_date,
|
||||
CASE
|
||||
WHEN http_user_agent LIKE '%bot%'
|
||||
OR http_user_agent LIKE '%spider%'
|
||||
OR http_user_agent LIKE '%crawl%'
|
||||
OR http_user_agent LIKE '%slurp%'
|
||||
OR http_user_agent LIKE '%bingpreview%'
|
||||
OR http_user_agent LIKE '%facebookexternalhit%'
|
||||
OR http_user_agent LIKE '%twitterbot%' THEN 'spider'
|
||||
WHEN http_user_agent IS NULL OR http_user_agent = '' OR http_user_agent = '-' THEN 'unknown'
|
||||
WHEN http_user_agent LIKE '%Mozilla%'
|
||||
OR http_user_agent LIKE '%Chrome%'
|
||||
OR http_user_agent LIKE '%Safari%'
|
||||
OR http_user_agent LIKE '%Firefox%'
|
||||
OR http_user_agent LIKE '%Edg%'
|
||||
OR http_user_agent LIKE '%Opera%'
|
||||
OR http_user_agent LIKE '%MSIE%'
|
||||
OR http_user_agent LIKE '%Trident%' THEN 'browser'
|
||||
ELSE 'unknown'
|
||||
END AS ua_type,
|
||||
CASE
|
||||
WHEN http_user_agent LIKE '%Googlebot%' THEN 'Googlebot'
|
||||
WHEN http_user_agent LIKE '%Baiduspider%' THEN 'Baiduspider'
|
||||
WHEN http_user_agent LIKE '%bingbot%' THEN 'Bingbot'
|
||||
WHEN http_user_agent LIKE '%DuckDuckBot%' THEN 'DuckDuckBot'
|
||||
WHEN http_user_agent LIKE '%YandexBot%' THEN 'YandexBot'
|
||||
WHEN http_user_agent LIKE '%Edg/%' THEN 'Microsoft Edge'
|
||||
WHEN http_user_agent LIKE '%OPR/%' OR http_user_agent LIKE '%Opera%' THEN 'Opera'
|
||||
WHEN http_user_agent LIKE '%Firefox/%' THEN 'Firefox'
|
||||
WHEN http_user_agent LIKE '%Chrome/%' THEN 'Chrome'
|
||||
WHEN http_user_agent LIKE '%Safari/%' THEN 'Safari'
|
||||
WHEN http_user_agent LIKE '%MSIE%' OR http_user_agent LIKE '%Trident%' THEN 'Internet Explorer'
|
||||
WHEN http_user_agent IS NULL OR http_user_agent = '' OR http_user_agent = '-' THEN 'Unknown'
|
||||
ELSE 'Other'
|
||||
END AS ua_name,
|
||||
COUNT(*) AS pv,
|
||||
COUNT(DISTINCT remote_addr) AS uv
|
||||
FROM `{$rawTable}`
|
||||
WHERE time_local >= :day_start AND time_local < :day_end {$nodeWhere}
|
||||
GROUP BY ua_type, ua_name
|
||||
ORDER BY pv DESC
|
||||
{$uaLimit}";
|
||||
$statUaRows = (int) Db::execute($uaSql, $dayParams);
|
||||
|
||||
return [
|
||||
'stat_hour' => $statHourRows,
|
||||
'stat_url' => $statUrlRows,
|
||||
'stat_referer' => $statRefererRows,
|
||||
'stat_ua' => $statUaRows,
|
||||
];
|
||||
}
|
||||
|
||||
/**
|
||||
* 将 int 日期 YYYYMMDD 格式化为 "YYYY-MM-DD".
|
||||
*/
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
namespace base\common\service;
|
||||
|
||||
use app\admin\model\NginxLogPosition;
|
||||
use app\common\service\HostService;
|
||||
|
||||
/**
|
||||
* Nginx 日志增量读取 service(Base 层).
|
||||
@@ -10,12 +11,17 @@ use app\admin\model\NginxLogPosition;
|
||||
* 职责:
|
||||
* - 按文件 offset 增量读取日志行(Generator,不读整个文件到内存)
|
||||
* - 检测日志轮转(inode 变化 OR filesize < offset,覆盖 rename+rebuild 与 copytruncate 两种模式)
|
||||
* - 将读取进度持久化到 NginxLogPosition 表(upsert)
|
||||
* - 将读取进度持久化到 NginxLogPosition 表(upsert,按 node_id 隔离)
|
||||
*
|
||||
* 多节点适配:position 表按 (node_id, file_path) 唯一定位读取进度。
|
||||
* 旧数据(node_id 为空)通过懒加载接管(见 getLastPosition):首个查询到它的本节点原子 UPDATE node_id,
|
||||
* 避免 migration 全表回填(多节点共享 DB 时 migration 只跑一次,只有懒加载能确保每节点各自接管)。
|
||||
*
|
||||
* 不负责:解析(T5 Parser)、聚合(T7 Aggregator)。
|
||||
*
|
||||
* 依赖倒置:内部 model 调用走 `app\admin\model\NginxLogPosition`(App 入口类),
|
||||
* 使用者可在 app/ 重写该 model 拦截行为;本类不直接 `new NginxLogPosition`(静态调用满足多态)。
|
||||
* 节点身份走 `app\common\service\HostService`(App 入口类)。
|
||||
*/
|
||||
class NginxLogReaderServiceBase
|
||||
{
|
||||
@@ -159,15 +165,27 @@ class NginxLogReaderServiceBase
|
||||
/**
|
||||
* 查询文件的上次读取位置.
|
||||
*
|
||||
* 多节点适配:where 条件带 node_id 隔离各节点进度。
|
||||
* 懒加载接管:多节点共享 DB 时 migration 只跑一次,旧记录 node_id 为空,
|
||||
* 首个查询到它的本节点原子 UPDATE node_id 接管,保证每节点各自有独立 position 记录。
|
||||
*
|
||||
* @param string $filePath 日志文件绝对路径
|
||||
*
|
||||
* @return array|null 命中时返回 [inode, offset, last_line_hash, parse_fail_count, parse_fail_samples],无记录返回 null
|
||||
*/
|
||||
public function getLastPosition(string $filePath): ?array
|
||||
{
|
||||
$row = NginxLogPosition::where('file_path', $filePath)->find();
|
||||
$nodeId = HostService::getNodeId();
|
||||
|
||||
$row = NginxLogPosition::where('file_path', $filePath)
|
||||
->where('node_id', $nodeId)
|
||||
->find();
|
||||
if ($row === null) {
|
||||
return null;
|
||||
// 懒加载接管:旧记录 node_id 为空,原子抢占接管到当前节点
|
||||
$row = $this->lazyAdoptPosition($filePath, $nodeId);
|
||||
if ($row === null) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
return [
|
||||
@@ -179,6 +197,39 @@ class NginxLogReaderServiceBase
|
||||
];
|
||||
}
|
||||
|
||||
/**
|
||||
* 懒加载接管:把旧的无 node_id 记录原子接管到当前节点.
|
||||
*
|
||||
* 多节点共享 DB 时 migration 只跑一次,无法为每节点预置 position 记录。
|
||||
* 旧记录(node_id='')由首个查询到它的本节点通过原子条件 UPDATE 接管:
|
||||
* UPDATE ... SET node_id=current WHERE node_id='' AND file_path=X
|
||||
* 并发安全:WHERE node_id='' 保证只有一个节点能 affected=1,其他节点查询时已被接管。
|
||||
*
|
||||
* @return \think\Model|null 接管成功返回接管后的记录,无旧记录或被其他节点抢占返回 null
|
||||
*/
|
||||
protected function lazyAdoptPosition(string $filePath, string $nodeId)
|
||||
{
|
||||
$legacy = NginxLogPosition::where('file_path', $filePath)
|
||||
->where('node_id', '')
|
||||
->find();
|
||||
if ($legacy === null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// 原子条件 UPDATE:只有 node_id='' 时才能被接管(防并发多节点同时接管同一记录)
|
||||
$affected = NginxLogPosition::where('id', $legacy->getData('id'))
|
||||
->where('node_id', '')
|
||||
->update(['node_id' => $nodeId]);
|
||||
if ($affected === 0) {
|
||||
// 被其他节点抢先接管,本节点放弃(下次 savePosition 会创建新记录)
|
||||
return null;
|
||||
}
|
||||
|
||||
return NginxLogPosition::where('file_path', $filePath)
|
||||
->where('node_id', $nodeId)
|
||||
->find();
|
||||
}
|
||||
|
||||
/**
|
||||
* 保存读取位置(upsert).
|
||||
*
|
||||
@@ -196,8 +247,11 @@ class NginxLogReaderServiceBase
|
||||
public function savePosition(string $filePath, int $inode, int $offset, ?string $lastLineHash, int $failCount = 0, ?string $failSamples = null): void
|
||||
{
|
||||
$now = time();
|
||||
$nodeId = HostService::getNodeId();
|
||||
|
||||
$existing = NginxLogPosition::where('file_path', $filePath)->find();
|
||||
$existing = NginxLogPosition::where('file_path', $filePath)
|
||||
->where('node_id', $nodeId)
|
||||
->find();
|
||||
if ($existing !== null) {
|
||||
$existing->save([
|
||||
'inode' => $inode,
|
||||
@@ -212,6 +266,7 @@ class NginxLogReaderServiceBase
|
||||
}
|
||||
|
||||
NginxLogPosition::create([
|
||||
'node_id' => $nodeId,
|
||||
'file_path' => $filePath,
|
||||
'inode' => $inode,
|
||||
'offset' => $offset,
|
||||
|
||||
Reference in New Issue
Block a user