Update webstats_common.lua

This commit is contained in:
dami 2026-07-04 22:43:20 +08:00
parent 77a0d3af36
commit 01ac598e39
1 changed files with 239 additions and 219 deletions

View File

@ -377,59 +377,67 @@ end
function _M.cronPre(self)
self:lock_working('cron_init_stat')
-- 获取当前小时和下一小时的存储键
local time_key = self:get_store_key()
local time_key_next = self:get_store_key_with_time(ngx.time() + 3600)
local ok, err = pcall(function()
-- 获取当前小时和下一小时的存储键
local time_key = self:get_store_key()
local time_key_next = self:get_store_key_with_time(ngx.time() + 3600)
-- 需要预创建的统计表
local wc_stat = {'request_stat', 'client_stat', 'spider_stat'}
-- 需要预创建的统计表
local wc_stat = {'request_stat', 'client_stat', 'spider_stat'}
-- 遍历所有站点
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
-- 遍历所有站点
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
-- 安全地初始化数据库连接
local ok, db = pcall(function() return self:initDB(input_sn) end)
if not ok or not db then
self:D("cronPre initDB failed for " .. tostring(input_sn) .. ": " .. tostring(db))
self:unlock_working('cron_init_stat')
return false
end
-- 开启事务
db:exec([[BEGIN TRANSACTION]])
local success = true
-- 预创建当前小时和下一小时的统计记录
for _, ws_v in ipairs(wc_stat) do
if not self:_update_stat_pre(db, ws_v, time_key) then
success = false
break
-- 安全地初始化数据库连接
local db_ok, db = pcall(function() return self:initDB(input_sn) end)
if not db_ok or not db then
self:D("cronPre initDB failed for " .. tostring(input_sn) .. ": " .. tostring(db))
return false
end
if not self:_update_stat_pre(db, ws_v, time_key_next) then
success = false
break
-- 开启事务
db:exec([[BEGIN TRANSACTION]])
local success = true
-- 预创建当前小时和下一小时的统计记录
for _, ws_v in ipairs(wc_stat) do
if not self:_update_stat_pre(db, ws_v, time_key) then
success = false
break
end
if not self:_update_stat_pre(db, ws_v, time_key_next) then
success = false
break
end
end
-- 提交或回滚事务
if success then
pcall(function() db:execute([[COMMIT]]) end)
else
pcall(function() db:exec([[ROLLBACK]]) end)
end
-- 安全关闭数据库
pcall(function() db:close() end)
if not success then
return false
end
end
-- 提交或回滚事务
if success then
pcall(function() db:execute([[COMMIT]]) end)
else
pcall(function() db:exec([[ROLLBACK]]) end)
end
-- 安全关闭数据库
pcall(function() db:close() end)
if not success then
self:unlock_working('cron_init_stat')
return false
end
end
return true
end)
self:unlock_working('cron_init_stat')
return true
if not ok then
self:D("cronPre error: " .. tostring(err))
return false
end
return err
end
--[[
@ -484,212 +492,224 @@ function _M.cron(self)
local time_key = self:get_store_key()
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
if self:is_migrating(input_sn) then
self:unlock_working(cron_key)
return true
end
if self:is_working('cron_init_stat') then
self:unlock_working(cron_key)
return true
end
local ok, db = pcall(function() return self:initDB(input_sn) end)
if not ok or not db then
self:D("initDB failed for " .. input_sn .. ": " .. tostring(db))
self:unlock_working(cron_key)
return true
end
stat_fields[input_sn] = {}
dbs[input_sn] = db
self:clean_stats(db, input_sn)
local ok, stmt = pcall(function()
return db:prepare[[INSERT INTO web_logs(
time, ip, domain, server_name, method, status_code, uri, body_length,
referer, user_agent, protocol, request_time, is_spider, request_headers, ip_list, client_port)
VALUES(:time, :ip, :domain, :server_name, :method, :status_code, :uri,
:body_length, :referer, :user_agent, :protocol, :request_time, :is_spider,
:request_headers, :ip_list, :client_port)]]
end)
if ok and stmt then
stmts[input_sn] = stmt
db:exec([[BEGIN TRANSACTION]])
else
-- self:D("prepare statement failed for " .. input_sn .. ": " .. tostring(stmt))
if db and db:isopen() then
db:close()
local function cleanup_and_unlock()
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
local stmt = stmts[input_sn]
if stmt then
pcall(function() stmt:finalize() end)
end
local db = dbs[input_sn]
if db and db:isopen() then
pcall(function() db:close() end)
end
dbs[input_sn] = nil
self:unlock_working(cron_key)
return true
end
self:unlock_working(cron_key)
end
local has_error = false
for i = 1, llen do
local data = ngx.shared.mw_total:lpop(total_key)
if not data then
break
local ok, err = pcall(function()
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
if self:is_migrating(input_sn) then
return
end
if self:is_working('cron_init_stat') then
return
end
local db_ok, db = pcall(function() return self:initDB(input_sn) end)
if not db_ok or not db then
self:D("initDB failed for " .. input_sn .. ": " .. tostring(db))
return
end
stat_fields[input_sn] = {}
dbs[input_sn] = db
self:clean_stats(db, input_sn)
local stmt_ok, stmt = pcall(function()
return db:prepare[[INSERT INTO web_logs(
time, ip, domain, server_name, method, status_code, uri, body_length,
referer, user_agent, protocol, request_time, is_spider, request_headers, ip_list, client_port)
VALUES(:time, :ip, :domain, :server_name, :method, :status_code, :uri,
:body_length, :referer, :user_agent, :protocol, :request_time, :is_spider,
:request_headers, :ip_list, :client_port)]]
end)
if stmt_ok and stmt then
stmts[input_sn] = stmt
db:exec([[BEGIN TRANSACTION]])
else
if db and db:isopen() then
db:close()
end
dbs[input_sn] = nil
return
end
end
local ok, info = pcall(function() return json.decode(data) end)
if not ok then
self:D("json decode failed: " .. tostring(info))
ngx.shared.mw_total:rpush(total_key, data)
has_error = true
break
end
for i = 1, llen do
local data = ngx.shared.mw_total:lpop(total_key)
if not data then
break
end
local input_sn = info['server_name']
local db = dbs[input_sn]
if not db then
ngx.shared.mw_total:rpush(total_key, data)
has_error = true
break
end
local decode_ok, info = pcall(function() return json.decode(data) end)
if not decode_ok then
self:D("json decode failed: " .. tostring(info))
ngx.shared.mw_total:rpush(total_key, data)
return
end
local stmt = stmts[input_sn]
if not stmt then
ngx.shared.mw_total:rpush(total_key, data)
has_error = true
break
end
local input_sn = info['server_name']
local db = dbs[input_sn]
if not db then
ngx.shared.mw_total:rpush(total_key, data)
return
end
local insert_ok = self:store_logs_line(db, stmt, input_sn, info)
if not insert_ok then
ngx.shared.mw_total:rpush(total_key, data)
rollback_sites[input_sn] = true
has_error = true
break
end
local stmt = stmts[input_sn]
if not stmt then
ngx.shared.mw_total:rpush(total_key, data)
return
end
local log_kv = info["log_kv"]
local excluded = log_kv['excluded']
local stat_tmp_fields = info['stat_fields']
local insert_ok = self:store_logs_line(db, stmt, input_sn, info)
if not insert_ok then
ngx.shared.mw_total:rpush(total_key, data)
rollback_sites[input_sn] = true
return
end
local stat_fields_is = stat_fields[input_sn]
for stf_k, stf_v in pairs(stat_tmp_fields) do
if excluded then
if stf_k == "spider_stat_fields" or stf_k == "client_stat_fields" then
break
local log_kv = info["log_kv"]
local excluded = log_kv['excluded']
local stat_tmp_fields = info['stat_fields']
local stat_fields_is = stat_fields[input_sn]
for stf_k, stf_v in pairs(stat_tmp_fields) do
if excluded then
if stf_k == "spider_stat_fields" or stf_k == "client_stat_fields" then
break
end
end
local stf_is = stat_fields_is[stf_k]
if not stf_is then
stf_is = {}
stat_fields_is[stf_k] = stf_is
end
for sv_k, sv_v in pairs(stf_v) do
stf_is[sv_k] = (stf_is[sv_k] or 0) + sv_v
end
end
local stf_is = stat_fields_is[stf_k]
if not stf_is then
stf_is = {}
stat_fields_is[stf_k] = stf_is
end
for sv_k, sv_v in pairs(stf_v) do
stf_is[sv_k] = (stf_is[sv_k] or 0) + sv_v
if not excluded then
local ip = log_kv['ip']
local body_length = log_kv["body_length"]
local ip_stats_sn = ip_stats[input_sn]
if not ip_stats_sn then
ip_stats_sn = {}
ip_stats[input_sn] = ip_stats_sn
end
local ip_stat = ip_stats_sn[ip]
if not ip_stat then
ip_stats_sn[ip] = {ip_num = 1, body_length = body_length}
else
ip_stat.ip_num = ip_stat.ip_num + 1
ip_stat.body_length = ip_stat.body_length + body_length
end
local url_stats_sn = url_stats[input_sn]
if not url_stats_sn then
url_stats_sn = {}
url_stats[input_sn] = url_stats_sn
end
local request_uri = log_kv["request_uri"]
local request_uri_md5 = ngx.md5(request_uri)
local url_stat = url_stats_sn[request_uri_md5]
if not url_stat then
url_stats_sn[request_uri_md5] = {url_num = 1, uri = request_uri, body_length = body_length}
else
url_stat.url_num = url_stat.url_num + 1
url_stat.body_length = url_stat.body_length + body_length
end
end
end
if not excluded then
local ip = log_kv['ip']
local body_length = log_kv["body_length"]
local save_day = config['global']["save_day"]
local now_date = os.date("*t")
local save_date_timestamp = os.time{year = now_date.year, month = now_date.month, day = now_date.day - save_day, hour = 0}
local delete_sql = "DELETE FROM web_logs WHERE time<" .. tostring(save_date_timestamp)
local ip_stats_sn = ip_stats[input_sn]
if not ip_stats_sn then
ip_stats_sn = {}
ip_stats[input_sn] = ip_stats_sn
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
local stmt = stmts[input_sn]
if stmt then
pcall(function() stmt:finalize() end)
end
local ip_stat = ip_stats_sn[ip]
if not ip_stat then
ip_stats_sn[ip] = {ip_num = 1, body_length = body_length}
else
ip_stat.ip_num = ip_stat.ip_num + 1
ip_stat.body_length = ip_stat.body_length + body_length
end
local stat_fields_is = stat_fields[input_sn]
local db = dbs[input_sn]
local url_stats_sn = url_stats[input_sn]
if not url_stats_sn then
url_stats_sn = {}
url_stats[input_sn] = url_stats_sn
end
local request_uri = log_kv["request_uri"]
local request_uri_md5 = ngx.md5(request_uri)
local url_stat = url_stats_sn[request_uri_md5]
if not url_stat then
url_stats_sn[request_uri_md5] = {url_num = 1, uri = request_uri, body_length = body_length}
else
url_stat.url_num = url_stat.url_num + 1
url_stat.body_length = url_stat.body_length + body_length
end
end
end
local save_day = config['global']["save_day"]
local now_date = os.date("*t")
local save_date_timestamp = os.time{year = now_date.year, month = now_date.month, day = now_date.day - save_day, hour = 0}
local delete_sql = "DELETE FROM web_logs WHERE time<" .. tostring(save_date_timestamp)
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
local stmt = stmts[input_sn]
if stmt then
pcall(function() stmt:finalize() end)
end
local stat_fields_is = stat_fields[input_sn]
local db = dbs[input_sn]
if db then
if rollback_sites[input_sn] then
pcall(function() db:exec([[ROLLBACK]]) end)
else
for sti_k, sti_v in pairs(stat_fields_is) do
local vkk = ""
for sv_k, sv_v in pairs(sti_v) do
vkk = vkk .. sv_k .. "=" .. sv_k .. "+" .. sv_v .. ","
end
if vkk ~= "" then
vkk = string.sub(vkk, 1, string.len(vkk) - 1)
if sti_k == 'request_stat_fields' then
self:update_stat(db, "request_stat", time_key, vkk)
elseif sti_k == 'client_stat_fields' then
self:update_stat(db, "client_stat", time_key, vkk)
elseif sti_k == 'spider_stat_fields' then
self:update_stat(db, "spider_stat", time_key, vkk)
if db then
if rollback_sites[input_sn] then
pcall(function() db:exec([[ROLLBACK]]) end)
else
for sti_k, sti_v in pairs(stat_fields_is) do
local vkk = ""
for sv_k, sv_v in pairs(sti_v) do
vkk = vkk .. sv_k .. "=" .. sv_k .. "+" .. sv_v .. ","
end
if vkk ~= "" then
vkk = string.sub(vkk, 1, string.len(vkk) - 1)
if sti_k == 'request_stat_fields' then
self:update_stat(db, "request_stat", time_key, vkk)
elseif sti_k == 'client_stat_fields' then
self:update_stat(db, "client_stat", time_key, vkk)
elseif sti_k == 'spider_stat_fields' then
self:update_stat(db, "spider_stat", time_key, vkk)
end
end
end
end
local local_ip_stats = ip_stats[input_sn]
if local_ip_stats then
for ip_addr, ip_val in pairs(local_ip_stats) do
self:update_statistics_ip(db, ip_addr, ip_val.ip_num, ip_val.body_length)
local local_ip_stats = ip_stats[input_sn]
if local_ip_stats then
for ip_addr, ip_val in pairs(local_ip_stats) do
self:update_statistics_ip(db, ip_addr, ip_val.ip_num, ip_val.body_length)
end
end
end
local local_url_stats = url_stats[input_sn]
if local_url_stats then
for url_md5, url_val in pairs(local_url_stats) do
self:update_statistics_uri(db, url_val.uri, url_md5, url_val.url_num, url_val.body_length)
local local_url_stats = url_stats[input_sn]
if local_url_stats then
for url_md5, url_val in pairs(local_url_stats) do
self:update_statistics_uri(db, url_val.uri, url_md5, url_val.url_num, url_val.body_length)
end
end
db:exec(delete_sql)
pcall(function() db:execute([[COMMIT]]) end)
end
end
db:exec(delete_sql)
pcall(function() db:execute([[COMMIT]]) end)
if db and db:isopen() then
pcall(function() db:close() end)
end
end
end)
if db and db:isopen() then
pcall(function() db:close() end)
end
if not ok then
self:D("cron process error: " .. tostring(err))
cleanup_and_unlock()
return true
end
self:unlock_working(cron_key)
@ -943,7 +963,7 @@ end
function _M.lock_working(self, sign)
local working_key = sign.."_working"
cache:set(working_key, true, 60)
cache:set(working_key, true, 10)
end
function _M.unlock_working(self, sign)