This commit is contained in:
dami
2026-07-06 19:24:36 +08:00
parent 58ab48ac15
commit f9be534101
2 changed files with 145 additions and 149 deletions
+132 -133
View File
@@ -397,41 +397,43 @@ function _M.cronPre(self)
for _, site_v in ipairs(sites) do
local input_sn = site_v["name"]
-- 安全地初始化数据库连接
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
-- 开启事务
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
if not self:is_migrating(input_sn) then
-- 安全地初始化数据库连接
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
end
-- 提交或回滚事务
if success then
pcall(function() db:execute([[COMMIT]]) end)
else
pcall(function() db:exec([[ROLLBACK]]) end)
end
-- 提交或回滚事务
if success then
pcall(function() db:execute([[COMMIT]]) end)
else
pcall(function() db:exec([[ROLLBACK]]) end)
end
-- 安全关闭数据库
pcall(function() db:close() end)
-- 安全关闭数据库
pcall(function() db:close() end)
if not success then
return false
if not success then
return false
end
end
end
@@ -515,45 +517,45 @@ function _M.cron(self)
end
local ok, err = pcall(function()
if self:is_working('cron_init_stat') then
return
end
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, scheme, 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, :scheme, :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()
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, scheme, 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, :scheme, :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
dbs[input_sn] = nil
return
end
end
@@ -571,82 +573,79 @@ function _M.cron(self)
end
local input_sn = info['server_name']
local db = dbs[input_sn]
if not db then
if self:is_migrating(input_sn) then
ngx.shared.mw_total:rpush(total_key, data)
return
end
else
local db = dbs[input_sn]
local stmt = stmts[input_sn]
if not db or not stmt then
ngx.shared.mw_total:rpush(total_key, data)
else
local insert_ok = self:store_logs_line(db, stmt, input_sn, info)
if not insert_ok then
self:D("store_logs_line failed for " .. input_sn .. ": " .. tostring(insert_ok))
ngx.shared.mw_total:rpush(total_key, data)
rollback_sites[input_sn] = true
else
local log_kv = info["log_kv"]
local excluded = log_kv['excluded']
local stat_tmp_fields = info['stat_fields']
local stmt = stmts[input_sn]
if not stmt then
ngx.shared.mw_total:rpush(total_key, data)
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
end
end
local insert_ok = self:store_logs_line(db, stmt, input_sn, info)
if not insert_ok then
self:D("store_logs_line failed for " .. input_sn .. ": " .. tostring(insert_ok))
ngx.shared.mw_total:rpush(total_key, data)
rollback_sites[input_sn] = true
return
end
local stf_is = stat_fields_is[stf_k]
if not stf_is then
stf_is = {}
stat_fields_is[stf_k] = stf_is
end
local log_kv = info["log_kv"]
local excluded = log_kv['excluded']
local stat_tmp_fields = info['stat_fields']
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 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
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
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
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
+13 -16
View File
@@ -88,7 +88,11 @@ def migrateSiteHotLogs(site_name, query_date):
print(f"[{site_name}] 正在迁移中,跳过")
return mw.returnMsg(True, f"{site_name} is migrating, skip")
# 1. 备份到临时文件(使用文件直接拷贝)
todayTime = time.strftime('%Y-%m-%d 00:00:00', time.localtime())
todayUt = int(time.mktime(time.strptime(
todayTime, "%Y-%m-%d %H:%M:%S")))
# 1. 备份到临时文件(copy 期间短暂互斥,完成后立即解除)
try:
import shutil
print(f"[{site_name}] 备份 {hot_db} -> {hot_db_tmp} ...")
@@ -98,11 +102,12 @@ def migrateSiteHotLogs(site_name, query_date):
if not os.path.exists(hot_db_tmp):
return mw.returnMsg(False, f"{site_name} migrating fail, copy tmp file!")
except Exception as e:
return mw.returnMsg(False, f"{site_name} migrating fail: {e}")
finally:
if os.path.exists(migrating_flag):
os.remove(migrating_flag)
return mw.returnMsg(False, f"{site_name} migrating fail: {e}")
# 2. 从临时备份中迁移热日志数据到历史日志(批量插入)
# 2. 从临时备份中迁移热日志数据到历史日志(读 logs_tmp,不阻塞 live 写入)
try:
logs_conn = pSqliteDb('web_log', site_name, 'logs_tmp')
history_logs_conn = pSqliteDb('web_log', site_name, 'history_logs')
@@ -113,10 +118,6 @@ def migrateSiteHotLogs(site_name, query_date):
columns_str = ",".join(_columns)
placeholders = ",".join(["?"] * len(_columns))
todayTime = time.strftime('%Y-%m-%d 00:00:00', time.localtime())
todayUt = int(time.mktime(time.strptime(
todayTime, "%Y-%m-%d %H:%M:%S")))
logs_sql = f"select {columns_str} from web_logs where time<{todayUt}"
selector = logs_conn.originExecute(logs_sql)
@@ -158,18 +159,15 @@ def migrateSiteHotLogs(site_name, query_date):
history_logs_conn.commit()
except Exception as e:
if site_name:
print(f"[{site_name}] logs to history error: {e}")
else:
print(f"logs to history error: {e}")
if os.path.exists(hot_db_tmp):
os.remove(hot_db_tmp)
print(f"[{site_name}] logs to history error: {e}")
return mw.returnMsg(False, f"{site_name} logs migrate error: {e}")
# 3. 删除已迁移的数据并清理统计(批量删除)
# 3. 删除已迁移的数据并清理统计(仅 DELETE/VACUUM 期间互斥)
try:
if os.path.exists(migrating_flag):
os.remove(migrating_flag)
mw.writeFile(migrating_flag, "yes")
time.sleep(0.5)
hot_db_conn = pSqliteDb('web_logs', site_name)
@@ -242,7 +240,6 @@ def migrateHotLogs(query_date="today"):
print(f"\n迁移完成! 成功: {success_count}, 失败: {fail_count}")
mw.opWeb('restart')
return mw.returnMsg(True, f"logs migrate ok, success: {success_count}, fail: {fail_count}")
except BlockingIOError: