mirror of
https://github.com/midoks/mdserver-web.git
synced 2026-10-10 00:39:25 +08:00
@@ -0,0 +1,209 @@
|
||||
# coding:utf-8
|
||||
|
||||
import sys
|
||||
import io
|
||||
import os
|
||||
import time
|
||||
import json
|
||||
|
||||
web_dir = os.getcwd() + "/web"
|
||||
if os.path.exists(web_dir):
|
||||
sys.path.append(web_dir)
|
||||
os.chdir(web_dir)
|
||||
|
||||
import core.mw as mw
|
||||
from utils.crontab import crontab as MwCrontab
|
||||
|
||||
app_debug = False
|
||||
if mw.isAppleSystem():
|
||||
app_debug = True
|
||||
|
||||
|
||||
def getPluginName():
|
||||
return 'webstats'
|
||||
|
||||
|
||||
def getPluginDir():
|
||||
return mw.getPluginDir() + '/' + getPluginName()
|
||||
|
||||
|
||||
def getServerDir():
|
||||
return mw.getServerDir() + '/' + getPluginName()
|
||||
|
||||
|
||||
def getConf():
|
||||
conf = getServerDir() + "/lua/config.json"
|
||||
return conf
|
||||
|
||||
|
||||
def getGlobalConf():
|
||||
conf = getConf()
|
||||
content = mw.readFile(conf)
|
||||
result = json.loads(content)
|
||||
return result
|
||||
|
||||
|
||||
def pSqliteDb(dbname='web_logs', site_name='unset', fn="logs"):
|
||||
|
||||
db_dir = getServerDir() + '/logs/' + site_name
|
||||
if not os.path.exists(db_dir):
|
||||
mw.execShell('mkdir -p ' + db_dir)
|
||||
|
||||
name = fn
|
||||
file = db_dir + '/' + name + '.db'
|
||||
|
||||
if not os.path.exists(file):
|
||||
conn = mw.M(dbname).dbPos(db_dir, name)
|
||||
sql = mw.readFile(getPluginDir() + '/conf/init.sql')
|
||||
sql_list = sql.split(';')
|
||||
for index in range(len(sql_list)):
|
||||
conn.execute(sql_list[index], ())
|
||||
else:
|
||||
conn = mw.M(dbname).dbPos(db_dir, name)
|
||||
|
||||
conn.execute("PRAGMA synchronous = 0", ())
|
||||
conn.execute("PRAGMA page_size = 4096", ())
|
||||
conn.execute("PRAGMA journal_mode = wal", ())
|
||||
|
||||
conn.autoTextFactory()
|
||||
|
||||
# conn.text_factory = lambda x: str(x, encoding="utf-8", errors='ignore')
|
||||
# conn.text_factory = lambda x: unicode(x, "utf-8", "ignore")
|
||||
return conn
|
||||
|
||||
|
||||
def migrateSiteHotLogs(site_name, query_date):
|
||||
print(site_name, query_date)
|
||||
|
||||
migrating_flag = getServerDir() + "/logs/%s/migrating" % site_name
|
||||
hot_db = getServerDir() + "/logs/%s/logs.db" % site_name
|
||||
hot_db_tmp = getServerDir() + "/logs/%s/logs_tmp.db" % site_name
|
||||
history_logs_db = getServerDir() + "/logs/%s/history_logs.db" % site_name
|
||||
# 1. copy to tmp file
|
||||
try:
|
||||
import shutil
|
||||
print("coping {} to {} ...".format(hot_db, hot_db_tmp))
|
||||
mw.writeFile(migrating_flag, "yes")
|
||||
time.sleep(3)
|
||||
shutil.copy(hot_db, hot_db_tmp)
|
||||
if not os.path.exists(hot_db_tmp):
|
||||
return mw.returnMsg(False, "migrating fail, copy tmp file!")
|
||||
except:
|
||||
return mw.returnMsg(False, "{} migrating fail.".format(site_name))
|
||||
finally:
|
||||
if os.path.exists(migrating_flag):
|
||||
os.remove(migrating_flag)
|
||||
|
||||
# 2. 从临时备份中迁移热日志数据到历史日志
|
||||
print("begin tmp to hot log data ...")
|
||||
try:
|
||||
print("history file: {}".format(history_logs_db))
|
||||
logs_conn = pSqliteDb('web_log', site_name, 'logs_tmp')
|
||||
history_logs_conn = pSqliteDb('web_log', site_name, 'history_logs')
|
||||
|
||||
hot_db_columns = logs_conn.originExecute(
|
||||
"PRAGMA table_info([web_logs])")
|
||||
_columns = ",".join([c[1] for c in hot_db_columns if c[1] != "id"])
|
||||
query_start = 0
|
||||
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 = "select {} from web_logs where time<{}".format(
|
||||
_columns, todayUt)
|
||||
selector = logs_conn.originExecute(logs_sql)
|
||||
log = selector.fetchone()
|
||||
while log:
|
||||
params = ""
|
||||
for field in log:
|
||||
if params:
|
||||
params += ","
|
||||
if field is None:
|
||||
field = "\'\'"
|
||||
elif type(field) == str:
|
||||
field = "\'" + field.replace("\'", "\"") + "\'"
|
||||
params += str(field)
|
||||
insert_sql = "insert into web_logs(" + \
|
||||
_columns + ") values(" + params + ")"
|
||||
history_logs_conn.execute(insert_sql)
|
||||
log = selector.fetchone()
|
||||
|
||||
print("sorting historical data, this action takes a long time...")
|
||||
history_logs_conn.execute("VACUUM;")
|
||||
|
||||
gcfg = getGlobalConf()
|
||||
save_day = gcfg['global']["save_day"]
|
||||
print("delete historical data {} days ago...".format(save_day))
|
||||
time_now = time.localtime()
|
||||
save_timestamp = time.mktime((time_now.tm_year, time_now.tm_mon, time_now.tm_mday - save_day, 0, 0, 0, 0, 0, 0))
|
||||
delete_sql = "delete from web_logs where time <= {}".format(
|
||||
save_timestamp)
|
||||
print('delete history_logs')
|
||||
print(delete_sql)
|
||||
history_logs_conn.execute(delete_sql)
|
||||
history_logs_conn.commit()
|
||||
|
||||
# 3. delete merged data and clean up statistics
|
||||
print("delete merged thermal data...")
|
||||
mw.writeFile(migrating_flag, "yes")
|
||||
|
||||
hot_db_conn = pSqliteDb('web_logs', site_name)
|
||||
del_hot_log = "delete from web_logs where time<{}".format(todayUt)
|
||||
print(del_hot_log)
|
||||
r = hot_db_conn.execute(del_hot_log)
|
||||
print("delete:", r)
|
||||
print("deleting statistics over 180 days...")
|
||||
save_time_key = time.strftime(
|
||||
'%Y%m%d00', time.localtime(time.time() - 180 * 86400))
|
||||
|
||||
del_request_stat_sql = "delete from request_stat where time<={}".format(
|
||||
save_time_key)
|
||||
hot_db_conn.execute(del_request_stat_sql)
|
||||
|
||||
hot_db_conn.execute(
|
||||
"delete from spider_stat where time<={}".format(save_time_key))
|
||||
hot_db_conn.execute(
|
||||
"delete from client_stat where time<={}".format(save_time_key))
|
||||
hot_db_conn.execute(
|
||||
"delete from referer_stat where time<={}".format(save_time_key))
|
||||
hot_db_conn.commit()
|
||||
print("clean up the hot database...")
|
||||
hot_db_conn.execute("VACUUM;")
|
||||
hot_db_conn.commit()
|
||||
|
||||
if os.path.exists(migrating_flag):
|
||||
os.remove(migrating_flag)
|
||||
except Exception as e:
|
||||
if site_name:
|
||||
print("{} logs to history error:{}".format(site_name, e))
|
||||
else:
|
||||
print("logs to history error:{}".format(e))
|
||||
finally:
|
||||
if os.path.exists(hot_db_tmp):
|
||||
os.remove(hot_db_tmp)
|
||||
|
||||
print("{} logs migrate ok.".format(site_name))
|
||||
|
||||
if not mw.isAppleSystem():
|
||||
mw.execShell("chown -R www:www " + getServerDir())
|
||||
|
||||
mw.opWeb('restart')
|
||||
return mw.returnMsg(True, "{} logs migrate ok".format(site_name))
|
||||
|
||||
|
||||
def migrateHotLogs(query_date="today"):
|
||||
print("begin migrate hot logs")
|
||||
sites = mw.M('sites').field('name').order("add_time").select()
|
||||
|
||||
unset_site = {"name": "unset"}
|
||||
sites.append(unset_site)
|
||||
|
||||
# migrateSiteHotLogs('t1.cn', query_date)
|
||||
|
||||
for site_info in sites:
|
||||
# print(site_info['name'])
|
||||
site_name = site_info["name"]
|
||||
migrate_res = migrateSiteHotLogs(site_name, query_date)
|
||||
if not migrate_res["status"]:
|
||||
print(migrate_res["msg"])
|
||||
print("end migrate hot logs")
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,667 @@
|
||||
log_by_lua_block {
|
||||
|
||||
local cpath = "{$SERVER_APP}/lua/"
|
||||
if not package.cpath:find(cpath) then
|
||||
package.cpath = cpath .. "?.so;" .. package.cpath
|
||||
end
|
||||
if not package.path:find(cpath) then
|
||||
package.path = cpath .. "?.lua;" .. package.path
|
||||
end
|
||||
|
||||
local ver = '0.2.4'
|
||||
local debug_mode = true
|
||||
|
||||
local __C = require "webstats_common"
|
||||
local C = __C:getInstance()
|
||||
|
||||
-- cache start ---
|
||||
local cache = ngx.shared.mw_total
|
||||
local function cache_set(server_name, id ,key, val)
|
||||
local line_kv = "log_kv_"..server_name..'_'..id.."_"..key
|
||||
-- cache:set(line_kv, val)
|
||||
cache:set(line_kv, val)
|
||||
end
|
||||
|
||||
local function cache_clear(server_name, id, key)
|
||||
local line_kv = "log_kv_"..server_name..'_'..id.."_"..key
|
||||
cache:delete(line_kv)
|
||||
end
|
||||
|
||||
local function cache_get(server_name, id, key)
|
||||
local line_kv = "log_kv_"..server_name..'_'..id.."_"..key
|
||||
local value = cache:get(line_kv)
|
||||
return value
|
||||
end
|
||||
|
||||
|
||||
-- cache end ---
|
||||
|
||||
-- domain config is import
|
||||
local db = nil
|
||||
local json = require "cjson"
|
||||
local sqlite3 = require "lsqlite3"
|
||||
local config = require "webstats_config"
|
||||
local sites = require "webstats_sites"
|
||||
|
||||
-- string.gsub(C:get_sn(ngx.var.server_name),'_','.')
|
||||
local server_name = C:get_sn(ngx.var.server_name)
|
||||
|
||||
|
||||
C:setConfData(config, sites)
|
||||
|
||||
local auto_config = C:setInputSn(server_name)
|
||||
|
||||
local request_header = ngx.req.get_headers()
|
||||
local method = ngx.req.get_method()
|
||||
local excluded = false
|
||||
|
||||
local day = os.date("%d")
|
||||
local number_day = tonumber(day)
|
||||
local day_column = "day"..number_day
|
||||
local flow_column = "flow"..number_day
|
||||
local spider_column = "spider_flow"..number_day
|
||||
--- default common var end ---
|
||||
|
||||
local function init_var()
|
||||
return true
|
||||
end
|
||||
|
||||
--------------------- exclude_func start --------------------------
|
||||
local function load_global_exclude_ip()
|
||||
local load_key = "global_exclude_ip_load"
|
||||
-- update global exclude ip
|
||||
local global_exclude_ip = auto_config["exclude_ip"]
|
||||
if global_exclude_ip then
|
||||
for i, _ip in pairs(global_exclude_ip)
|
||||
do
|
||||
-- global
|
||||
-- D("set global exclude ip: ".._ip)
|
||||
if not cache:get("global_exclude_ip_".._ip) then
|
||||
cache:set("global_exclude_ip_".._ip, true)
|
||||
end
|
||||
end
|
||||
end
|
||||
-- set tag
|
||||
cache:set(load_key, true)
|
||||
end
|
||||
|
||||
local function load_exclude_ip(input_server_name)
|
||||
|
||||
local load_key = input_server_name .. "_exclude_ip_load"
|
||||
local site_config = config[input_server_name]
|
||||
|
||||
local site_exclude_ip = nil
|
||||
if site_config then
|
||||
site_exclude_ip = site_config["exclude_ip"]
|
||||
end
|
||||
|
||||
-- update server_name exclude ip
|
||||
if site_exclude_ip then
|
||||
for i, _ip in pairs(site_exclude_ip)
|
||||
do
|
||||
cache:set(input_server_name .. "_exclude_ip_".._ip, true)
|
||||
end
|
||||
end
|
||||
|
||||
-- set tag
|
||||
cache:set(load_key, true)
|
||||
return true
|
||||
end
|
||||
|
||||
local function filter_status()
|
||||
if not auto_config['exclude_status'] then return false end
|
||||
local the_status = tostring(ngx.status)
|
||||
for _,v in ipairs(auto_config['exclude_status'])
|
||||
do
|
||||
if the_status == v then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
local function exclude_extension()
|
||||
if not ngx.var.uri then return false end
|
||||
if not auto_config['exclude_extension'] then return false end
|
||||
for _,v in ipairs(auto_config['exclude_extension'])
|
||||
do
|
||||
if ngx.re.find(ngx.var.uri,"[.]"..v.."$",'ijo') then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
local function exclude_url()
|
||||
if not ngx.var.uri then return false end
|
||||
if not ngx.var.request_uri then return false end
|
||||
if not auto_config['exclude_url'] then return false end
|
||||
local the_uri = string.sub(ngx.var.request_uri, 2)
|
||||
local url_conf = auto_config["exclude_url"]
|
||||
for i,conf in ipairs(url_conf)
|
||||
do
|
||||
local mode = conf["mode"]
|
||||
local url = conf["url"]
|
||||
if mode == "regular" then
|
||||
if ngx.re.find(the_uri, url, "ijo") then
|
||||
return true
|
||||
end
|
||||
else
|
||||
if the_uri == url then
|
||||
return true
|
||||
end
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
local function exclude_ip(input_server_name, ip)
|
||||
-- 排除IP匹配,分网站单独的配置和全局配置两种方式
|
||||
local site_config = config[input_server_name]
|
||||
local server_exclude_ips = nil
|
||||
local check_server_exclude_ip = false
|
||||
if site_config then
|
||||
server_exclude_ips = site_config["exclude_ip"]
|
||||
if not server_exclude_ips then
|
||||
return false
|
||||
end
|
||||
for k, _ip in pairs(server_exclude_ips)
|
||||
do
|
||||
check_server_exclude_ip = true
|
||||
break
|
||||
end
|
||||
end
|
||||
-- D("server[" ..input_server_name.."]check exclude ip : "..tostring(check_server_exclude_ip))
|
||||
if check_server_exclude_ip then
|
||||
if cache:get(input_server_name .. "_exclude_ip_"..ip) then
|
||||
-- D("-Exclude server ip:"..ip)
|
||||
return true
|
||||
end
|
||||
else
|
||||
if cache:get("global_exclude_ip_"..ip) then
|
||||
-- D("*Excluded global ip:"..ip)
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
--------------------- exclude_func end ---------------------------
|
||||
|
||||
local function cache_logs_old(server_name)
|
||||
|
||||
-- make new id
|
||||
local new_id = C:get_last_id(server_name)
|
||||
|
||||
local excluded = false
|
||||
local ip = C:get_client_ip()
|
||||
excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(server_name, ip)
|
||||
|
||||
local ip_list = request_header["x-forwarded-for"]
|
||||
if ip and not ip_list then
|
||||
ip_list = ip
|
||||
end
|
||||
|
||||
local remote_addr = ngx.var.remote_addr
|
||||
if not string.find(ip_list, remote_addr) then
|
||||
if remote_addr then
|
||||
ip_list = ip_list .. "," .. remote_addr
|
||||
end
|
||||
end
|
||||
|
||||
-- local request_time = ngx.var.request_time
|
||||
local request_time = C:get_request_time()
|
||||
local client_port = ngx.var.remote_port
|
||||
local real_server_name = ngx.var.server_name
|
||||
local uri = ngx.var.uri
|
||||
local status_code = ngx.status
|
||||
local protocol = ngx.var.server_protocol
|
||||
local request_uri = ngx.var.request_uri
|
||||
local time_key = C:get_store_key()
|
||||
local method = method
|
||||
local body_length = C:get_length()
|
||||
local domain = C:get_domain()
|
||||
local referer = ngx.var.http_referer
|
||||
|
||||
local kv = {
|
||||
id=new_id,
|
||||
time_key=time_key,
|
||||
time=ngx.time(),
|
||||
ip=ip,
|
||||
domain=domain,
|
||||
server_name=server_name,
|
||||
real_server_name=real_server_name,
|
||||
method=method,
|
||||
status_code=status_code,
|
||||
uri=uri,
|
||||
request_uri=request_uri,
|
||||
body_length=body_length,
|
||||
referer=referer,
|
||||
user_agent=request_header['user-agent'],
|
||||
protocol=protocol,
|
||||
is_spider=0,
|
||||
request_time=request_time,
|
||||
excluded=excluded,
|
||||
request_headers='',
|
||||
ip_list=ip_list,
|
||||
client_port=client_port
|
||||
}
|
||||
|
||||
local request_stat_fields = "req=req+1,length=length+"..body_length
|
||||
local spider_stat_fields = "x"
|
||||
local client_stat_fields = "x"
|
||||
|
||||
if not excluded then
|
||||
|
||||
if status_code == 500 or (method=="POST" and config["record_post_args"] == true) or (status_code==403 and config["record_get_403_args"] == true) then
|
||||
local data = ""
|
||||
local ok, err = pcall(function() data = C:get_http_origin() end)
|
||||
if ok and not err then
|
||||
kv["request_headers"] = data
|
||||
end
|
||||
end
|
||||
|
||||
if ngx.re.find("500,501,502,503,504,505,506,507,509,510,400,401,402,403,404,405,406,407,408,409,410,411,412,413,414,415,416,417,418,421,422,423,424,425,426,449,451,499", tostring(status_code), "jo") then
|
||||
local field = "status_"..status_code
|
||||
request_stat_fields = request_stat_fields .. ","..field.."="..field.."+1"
|
||||
end
|
||||
|
||||
local lower_method = string.lower(method)
|
||||
if ngx.re.find("get,post,put,patch,delete", lower_method, "ijo") then
|
||||
local field = "http_"..lower_method
|
||||
request_stat_fields = request_stat_fields .. ","..field.."="..field.."+1"
|
||||
end
|
||||
|
||||
|
||||
local ipc = 0
|
||||
local pvc = 0
|
||||
local uvc = 0
|
||||
|
||||
local is_spider, request_spider, spider_index = C:match_spider(kv['user_agent'])
|
||||
if not is_spider then
|
||||
|
||||
client_stat_fields = C:match_client(kv['user_agent'])
|
||||
if not client_stat_fields or #client_stat_fields == 0 then
|
||||
client_stat_fields = request_stat_fields..",other=other+1"
|
||||
end
|
||||
|
||||
pvc, uvc = C:statistics_request(ip, is_spider,body_length)
|
||||
ipc = C:statistics_ipc(server_name,ip)
|
||||
else
|
||||
kv["is_spider"] = spider_index
|
||||
local field = "spider"
|
||||
spider_stat_fields = request_spider.."="..request_spider.."+"..1
|
||||
request_stat_fields = request_stat_fields .. ","..field.."="..field.."+"..1
|
||||
end
|
||||
|
||||
if ipc > 0 then
|
||||
request_stat_fields = request_stat_fields..",ip=ip+1"
|
||||
end
|
||||
if uvc > 0 then
|
||||
request_stat_fields = request_stat_fields..",uv=uv+1"
|
||||
end
|
||||
if pvc > 0 then
|
||||
request_stat_fields = request_stat_fields..",pv=pv+1"
|
||||
end
|
||||
end
|
||||
|
||||
local stat_fields = request_stat_fields..";"..client_stat_fields..";"..spider_stat_fields
|
||||
|
||||
cache_set(server_name, new_id, "stat_fields", stat_fields)
|
||||
cache_set(server_name, new_id, "log_kv", json.encode(kv))
|
||||
end
|
||||
|
||||
local function cache_logs(input_sn)
|
||||
|
||||
-- make new id
|
||||
local new_id = C:get_last_id(input_sn)
|
||||
|
||||
local excluded = false
|
||||
local ip = C:get_client_ip()
|
||||
excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(input_sn, ip)
|
||||
|
||||
local ip_list = request_header["x-forwarded-for"]
|
||||
if ip and not ip_list then
|
||||
ip_list = ip
|
||||
end
|
||||
|
||||
local remote_addr = ngx.var.remote_addr
|
||||
|
||||
-- 修复反向代理代过来的数据
|
||||
if "table" == type(ip_list) then
|
||||
ip_list = json.encode(ip_list)
|
||||
end
|
||||
|
||||
if not string.find(ip_list, remote_addr) then
|
||||
if remote_addr then
|
||||
ip_list = ip_list .. "," .. remote_addr
|
||||
end
|
||||
end
|
||||
|
||||
-- local request_time = ngx.var.request_time
|
||||
local request_time = C:get_request_time()
|
||||
local client_port = ngx.var.remote_port
|
||||
local real_server_name = input_sn
|
||||
local uri = tostring(ngx.var.uri)
|
||||
local status_code = ngx.status
|
||||
local protocol = ngx.var.server_protocol
|
||||
local request_uri = ngx.var.request_uri
|
||||
local time_key = C:get_store_key()
|
||||
local method = method
|
||||
local body_length = C:get_length()
|
||||
local domain = C:get_domain()
|
||||
local referer = ngx.var.http_referer
|
||||
|
||||
local kv = {
|
||||
id=new_id,
|
||||
time_key=time_key,
|
||||
time=ngx.time(),
|
||||
ip=ip,
|
||||
domain=domain,
|
||||
server_name=input_sn,
|
||||
real_server_name=real_server_name,
|
||||
method=method,
|
||||
status_code=status_code,
|
||||
uri=uri,
|
||||
request_uri=request_uri,
|
||||
body_length=body_length,
|
||||
referer=referer,
|
||||
user_agent=request_header['user-agent'],
|
||||
protocol=protocol,
|
||||
is_spider=0,
|
||||
request_time=request_time,
|
||||
excluded=excluded,
|
||||
request_headers='',
|
||||
ip_list=ip_list,
|
||||
client_port=client_port
|
||||
}
|
||||
|
||||
-- C:D(json.encode(kv))
|
||||
local request_stat_fields = {req=1,length=body_length}
|
||||
|
||||
local spider_stat_fields = {}
|
||||
local client_stat_fields = {}
|
||||
|
||||
if not excluded then
|
||||
if status_code == 500 or (method=="POST" and config['global']["record_post_args"] == true) or (status_code==403 and config['global']["record_get_403_args"] == true) then
|
||||
local data = ""
|
||||
local ok, err = pcall(function() data = C:get_http_origin() end)
|
||||
if ok and not err then
|
||||
kv["request_headers"] = data
|
||||
else
|
||||
C:D("debug request_headers error:"..tostring(err))
|
||||
end
|
||||
end
|
||||
|
||||
if ngx.re.find("500,501,502,503,504,505,506,507,509,510,400,401,402,403,404,405,406,407,408,409,410,411,412,413,414,415,416,417,418,421,422,423,424,425,426,449,451,499", tostring(status_code), "jo") then
|
||||
local field = "status_"..status_code
|
||||
request_stat_fields[field] = 1
|
||||
end
|
||||
|
||||
-- D("method:"..method)
|
||||
local lower_method = string.lower(method)
|
||||
if ngx.re.find("get,post,put,patch,delete", lower_method, "ijo") then
|
||||
local field = "http_"..lower_method
|
||||
request_stat_fields[field] = 1
|
||||
end
|
||||
|
||||
|
||||
local ipc = 0
|
||||
local pvc = 0
|
||||
local uvc = 0
|
||||
|
||||
local is_spider, request_spider, spider_index = C:match_spider(kv['user_agent'])
|
||||
if not is_spider then
|
||||
|
||||
client_stat_fields = C:match_client_arr(kv['user_agent'])
|
||||
pvc, uvc = C:statistics_request(ip, is_spider,body_length)
|
||||
ipc = C:statistics_ipc(input_sn,ip)
|
||||
else
|
||||
kv["is_spider"] = spider_index
|
||||
local field = "spider"
|
||||
spider_stat_fields[request_spider] = 1
|
||||
request_stat_fields[field] = 1
|
||||
end
|
||||
|
||||
if ipc > 0 then
|
||||
request_stat_fields["ip"] = 1
|
||||
end
|
||||
if uvc > 0 then
|
||||
request_stat_fields["uv"] = 1
|
||||
end
|
||||
if pvc > 0 then
|
||||
request_stat_fields["pv"] = 1
|
||||
end
|
||||
end
|
||||
|
||||
local stat_fields = {
|
||||
request_stat_fields=request_stat_fields,
|
||||
client_stat_fields=client_stat_fields,
|
||||
spider_stat_fields=spider_stat_fields,
|
||||
}
|
||||
|
||||
local data = {
|
||||
server_name=input_sn,
|
||||
stat_fields=stat_fields,
|
||||
log_kv=kv,
|
||||
}
|
||||
|
||||
local push_data = json.encode(data)
|
||||
-- C:D(json.encode(push_data))
|
||||
local key = C:getTotalKey()
|
||||
ngx.shared.mw_total:rpush(key, push_data)
|
||||
end
|
||||
|
||||
local function store_logs_line(db, stmt, input_server_name, lineno)
|
||||
local logvalue = cache_get(input_server_name, lineno, "log_kv")
|
||||
if not logvalue then return false end
|
||||
local logline = json.decode(logvalue)
|
||||
|
||||
local time = logline["time"]
|
||||
local id = logline["id"]
|
||||
local protocol = logline["protocol"]
|
||||
local client_port = logline["client_port"]
|
||||
local status_code = logline["status_code"]
|
||||
local uri = logline["uri"]
|
||||
local request_uri = logline["request_uri"]
|
||||
local method = logline["method"]
|
||||
local body_length = logline["body_length"]
|
||||
local referer = logline["referer"]
|
||||
local ip = logline["ip"]
|
||||
local ip_list = logline["ip_list"]
|
||||
local request_time = logline["request_time"]
|
||||
local is_spider = logline["is_spider"]
|
||||
local domain = logline["domain"]
|
||||
local server_name = logline["server_name"]
|
||||
local user_agent = logline["user_agent"]
|
||||
local request_headers = logline["request_headers"]
|
||||
local excluded = logline["excluded"]
|
||||
|
||||
local request_stat_fields = nil
|
||||
local client_stat_fields = nil
|
||||
local spider_stat_fields = nil
|
||||
local stat_fields = cache_get(input_server_name, id, "stat_fields")
|
||||
if stat_fields == nil then
|
||||
-- D("Log stat fields is nil.")
|
||||
-- D("Logdata:"..logvalue)
|
||||
else
|
||||
stat_fields = C:split(stat_fields, ";")
|
||||
request_stat_fields = stat_fields[1]
|
||||
client_stat_fields = stat_fields[2]
|
||||
spider_stat_fields = stat_fields[3]
|
||||
|
||||
if "x" == client_stat_fields then
|
||||
client_stat_fields = nil
|
||||
end
|
||||
|
||||
if "x" == spider_stat_fields then
|
||||
spider_stat_fields = nil
|
||||
end
|
||||
end
|
||||
|
||||
local time_key = logline["time_key"]
|
||||
if not excluded then
|
||||
stmt:bind_names{
|
||||
time=time,
|
||||
ip=ip,
|
||||
domain=domain,
|
||||
server_name=server_name,
|
||||
method=method,
|
||||
status_code=status_code,
|
||||
uri=request_uri,
|
||||
body_length=body_length,
|
||||
referer=referer,
|
||||
user_agent=user_agent,
|
||||
protocol=protocol,
|
||||
request_time=request_time,
|
||||
is_spider=is_spider,
|
||||
request_headers=request_headers,
|
||||
ip_list=ip_list,
|
||||
client_port=client_port,
|
||||
}
|
||||
|
||||
local res, err = stmt:step()
|
||||
if tostring(res) == "5" then
|
||||
-- D("step res:"..tostring(res))
|
||||
-- D("step err:"..tostring(err))
|
||||
-- D("the step database connection is busy, so it will be stored later.")
|
||||
return false
|
||||
end
|
||||
stmt:reset()
|
||||
-- D("store_logs_line ok")
|
||||
|
||||
C:update_stat( db, "client_stat", time_key, client_stat_fields)
|
||||
C:update_stat( db, "spider_stat", time_key, spider_stat_fields)
|
||||
-- C:D("stat ok")
|
||||
|
||||
-- only count non spider requests
|
||||
local ok, err = C:statistics_uri(db, request_uri, ngx.md5(request_uri), body_length)
|
||||
local ok, err = C:statistics_ip(db, ip, body_length)
|
||||
end
|
||||
|
||||
C:update_stat( db, "request_stat", time_key, request_stat_fields)
|
||||
return true
|
||||
end
|
||||
|
||||
local function store_logs(input_sn)
|
||||
if C:is_migrating(input_sn) == true then
|
||||
-- D("migrating...")
|
||||
return
|
||||
end
|
||||
|
||||
local last_insert_id_key = input_sn.."_last_id"
|
||||
local store_start_id_key = input_sn.."_store_start"
|
||||
local last_id = cache:get(last_insert_id_key)
|
||||
local store_start = cache:get(store_start_id_key)
|
||||
if store_start == nil then
|
||||
store_start = 1
|
||||
end
|
||||
local store_end = last_id
|
||||
if store_end == nil then
|
||||
store_end = 1
|
||||
end
|
||||
|
||||
local worker_id = ngx.worker.id()
|
||||
if C:is_working(input_sn) then
|
||||
-- D("other workers are being stored, please store later.")
|
||||
-- cache:delete(flush_data_key)
|
||||
return true
|
||||
end
|
||||
C:lock_working(input_sn)
|
||||
|
||||
local log_dir = "{$SERVER_APP}/logs"
|
||||
local db_path = log_dir .. '/' .. input_sn .. "/logs.db"
|
||||
local db, err = sqlite3.open(db_path)
|
||||
|
||||
if tostring(err) ~= 'nil' then
|
||||
C:D("sqlite3 open error:"..tostring(err))
|
||||
return true
|
||||
end
|
||||
|
||||
db:exec([[PRAGMA synchronous = 0]])
|
||||
db:exec([[PRAGMA page_size = 4096]])
|
||||
db:exec([[PRAGMA journal_mode = wal]])
|
||||
db:exec([[PRAGMA journal_size_limit = 1073741824]])
|
||||
|
||||
local stmt2 = nil
|
||||
if db ~= nil then
|
||||
stmt2 = 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 db == nil or stmt2 == nil then
|
||||
-- D("web data db error")
|
||||
-- cache:set(storing_key, false)
|
||||
if db and db:isopen() then
|
||||
db:close()
|
||||
end
|
||||
return true
|
||||
end
|
||||
|
||||
status, errorString = db:exec([[BEGIN TRANSACTION]])
|
||||
|
||||
C:clean_stats(db, input_sn)
|
||||
|
||||
if store_end >= store_start then
|
||||
for i=store_start, store_end, 1 do
|
||||
-- D("store_start:"..store_start..":store_end:".. store_end)
|
||||
if store_logs_line(db, stmt2, input_sn, i) then
|
||||
cache_clear(input_sn, i, "log_kv")
|
||||
cache_clear(input_sn, i, "stat_fields")
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
local res, err = stmt2:finalize()
|
||||
if tostring(res) == "5" then
|
||||
C:D("Finalize res:"..tostring(res)..",Finalize err:"..tostring(err))
|
||||
end
|
||||
|
||||
local now_date = os.date("*t")
|
||||
local save_day = config['global']["save_day"]
|
||||
local save_date_timestamp = os.time{year=now_date.year,
|
||||
month=now_date.month, day=now_date.day-save_day, hour=0}
|
||||
-- delete expire data
|
||||
db:exec("DELETE FROM web_logs WHERE time<"..tostring(save_date_timestamp))
|
||||
|
||||
local res, err = db:execute([[COMMIT]])
|
||||
if db and db:isopen() then
|
||||
db:close()
|
||||
end
|
||||
|
||||
cache:set(store_start_id_key, store_end+1)
|
||||
C:unlock_working(input_sn)
|
||||
end
|
||||
|
||||
local function run_app()
|
||||
-- C:D("------------ debug start ------------")
|
||||
init_var()
|
||||
|
||||
load_global_exclude_ip()
|
||||
load_exclude_ip(server_name)
|
||||
|
||||
cache_logs(server_name)
|
||||
-- C:cron()
|
||||
|
||||
-- cache_logs_old(server_name)
|
||||
-- store_logs(server_name)
|
||||
-- C:D("------------ debug end -------------")
|
||||
end
|
||||
|
||||
|
||||
local function run_app_ok()
|
||||
if not debug_mode then return run_app() end
|
||||
|
||||
local presult, err = pcall( function() run_app() end)
|
||||
if not presult then
|
||||
C:D("debug error on :"..tostring(err))
|
||||
return true
|
||||
end
|
||||
end
|
||||
|
||||
return run_app_ok()
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,5 +1,28 @@
|
||||
log_by_lua_block {
|
||||
--[[
|
||||
webstats_log.lua - WebStats 日志采集模块
|
||||
=========================================
|
||||
|
||||
本模块在 Nginx 的 log_by_lua_block 阶段执行,负责:
|
||||
- 采集请求日志数据
|
||||
- 过滤不需要统计的请求
|
||||
- 识别爬虫和客户端类型
|
||||
- 计算 PV/UV/IP 等指标
|
||||
- 将日志数据发送到共享内存队列
|
||||
|
||||
执行时机:每次请求完成后(log_by_lua_block 阶段)
|
||||
|
||||
主要数据流向:
|
||||
请求完成 -> log_by_lua_block -> 数据采集 -> 过滤处理 -> 爬虫/客户端识别 -> 统计计算 -> 入队
|
||||
|
||||
依赖模块:
|
||||
- cjson: JSON 编解码
|
||||
- webstats_common: 核心工具模块
|
||||
- webstats_config: 全局配置
|
||||
- webstats_sites: 站点配置
|
||||
]]
|
||||
|
||||
-- 添加 Lua 模块搜索路径
|
||||
local cpath = "{$SERVER_APP}/lua/"
|
||||
if not package.cpath:find(cpath) then
|
||||
package.cpath = cpath .. "?.so;" .. package.cpath
|
||||
@@ -8,85 +31,90 @@ log_by_lua_block {
|
||||
package.path = cpath .. "?.lua;" .. package.path
|
||||
end
|
||||
|
||||
local ver = '0.2.4'
|
||||
-- 调试模式开关
|
||||
local debug_mode = true
|
||||
|
||||
-- 引入核心模块(单例模式)
|
||||
local __C = require "webstats_common"
|
||||
local C = __C:getInstance()
|
||||
|
||||
-- cache start ---
|
||||
local cache = ngx.shared.mw_total
|
||||
local function cache_set(server_name, id ,key, val)
|
||||
local line_kv = "log_kv_"..server_name..'_'..id.."_"..key
|
||||
-- cache:set(line_kv, val)
|
||||
cache:set(line_kv, val)
|
||||
end
|
||||
|
||||
local function cache_clear(server_name, id, key)
|
||||
local line_kv = "log_kv_"..server_name..'_'..id.."_"..key
|
||||
cache:delete(line_kv)
|
||||
end
|
||||
|
||||
local function cache_get(server_name, id, key)
|
||||
local line_kv = "log_kv_"..server_name..'_'..id.."_"..key
|
||||
local value = cache:get(line_kv)
|
||||
return value
|
||||
end
|
||||
|
||||
|
||||
-- cache end ---
|
||||
|
||||
-- domain config is import
|
||||
local db = nil
|
||||
-- 依赖模块
|
||||
local json = require "cjson"
|
||||
local sqlite3 = require "lsqlite3"
|
||||
local config = require "webstats_config"
|
||||
local sites = require "webstats_sites"
|
||||
|
||||
-- string.gsub(C:get_sn(ngx.var.server_name),'_','.')
|
||||
-- 获取站点名称(通过域名匹配)
|
||||
local server_name = C:get_sn(ngx.var.server_name)
|
||||
|
||||
|
||||
-- 设置配置数据
|
||||
C:setConfData(config, sites)
|
||||
|
||||
-- 获取合并后的站点配置
|
||||
local auto_config = C:setInputSn(server_name)
|
||||
|
||||
-- 获取请求头部和方法
|
||||
local request_header = ngx.req.get_headers()
|
||||
local method = ngx.req.get_method()
|
||||
local excluded = false
|
||||
|
||||
local day = os.date("%d")
|
||||
local number_day = tonumber(day)
|
||||
local day_column = "day"..number_day
|
||||
local flow_column = "flow"..number_day
|
||||
local spider_column = "spider_flow"..number_day
|
||||
--- default common var end ---
|
||||
-- 共享内存字典实例(必须在使用前定义)
|
||||
local cache = ngx.shared.mw_total
|
||||
-- 日志队列键名
|
||||
local total_key = "log_kv_total"
|
||||
|
||||
local function init_var()
|
||||
return true
|
||||
end
|
||||
--[[
|
||||
状态码过滤表
|
||||
需要记录详细信息的状态码(4xx 和 5xx 错误)
|
||||
]]
|
||||
local status_codes_to_log = {
|
||||
["400"] = true, ["401"] = true, ["402"] = true, ["403"] = true, ["404"] = true,
|
||||
["405"] = true, ["406"] = true, ["407"] = true, ["408"] = true, ["409"] = true,
|
||||
["410"] = true, ["411"] = true, ["412"] = true, ["413"] = true, ["414"] = true,
|
||||
["415"] = true, ["416"] = true, ["417"] = true, ["418"] = true, ["421"] = true,
|
||||
["422"] = true, ["423"] = true, ["424"] = true, ["425"] = true, ["426"] = true,
|
||||
["449"] = true, ["451"] = true, ["499"] = true, ["500"] = true, ["501"] = true,
|
||||
["502"] = true, ["503"] = true, ["504"] = true, ["505"] = true, ["506"] = true,
|
||||
["507"] = true, ["509"] = true, ["510"] = true
|
||||
}
|
||||
|
||||
--------------------- exclude_func start --------------------------
|
||||
--[[
|
||||
HTTP 方法过滤表
|
||||
需要统计的 HTTP 方法
|
||||
]]
|
||||
local http_methods = {["get"] = true, ["post"] = true, ["put"] = true, ["patch"] = true, ["delete"] = true}
|
||||
|
||||
--[[
|
||||
排除函数模块
|
||||
============
|
||||
|
||||
以下函数用于判断请求是否应该被排除(不统计),包括:
|
||||
- IP 排除:全局和站点级别的排除 IP
|
||||
- 状态码排除:指定状态码的请求不统计
|
||||
- 扩展名排除:指定扩展名的文件不统计
|
||||
- URL 排除:指定 URL 路径不统计
|
||||
]]
|
||||
|
||||
--[[
|
||||
加载全局排除 IP 列表到缓存
|
||||
将配置中的全局排除 IP 存入共享内存,避免每次请求都读取配置
|
||||
]]
|
||||
local function load_global_exclude_ip()
|
||||
local load_key = "global_exclude_ip_load"
|
||||
-- update global exclude ip
|
||||
local global_exclude_ip = auto_config["exclude_ip"]
|
||||
if global_exclude_ip then
|
||||
for i, _ip in pairs(global_exclude_ip)
|
||||
do
|
||||
-- global
|
||||
-- D("set global exclude ip: ".._ip)
|
||||
if not cache:get("global_exclude_ip_".._ip) then
|
||||
cache:set("global_exclude_ip_".._ip, true)
|
||||
for _, _ip in pairs(global_exclude_ip) do
|
||||
if not cache:get("global_exclude_ip_" .. _ip) then
|
||||
cache:set("global_exclude_ip_" .. _ip, true)
|
||||
end
|
||||
end
|
||||
end
|
||||
-- set tag
|
||||
cache:set(load_key, true)
|
||||
end
|
||||
|
||||
--[[
|
||||
加载站点级排除 IP 列表到缓存
|
||||
@param input_server_name string 站点名称
|
||||
@return boolean 是否成功
|
||||
]]
|
||||
local function load_exclude_ip(input_server_name)
|
||||
|
||||
local load_key = input_server_name .. "_exclude_ip_load"
|
||||
local site_config = config[input_server_name]
|
||||
|
||||
@@ -95,24 +123,25 @@ log_by_lua_block {
|
||||
site_exclude_ip = site_config["exclude_ip"]
|
||||
end
|
||||
|
||||
-- update server_name exclude ip
|
||||
if site_exclude_ip then
|
||||
for i, _ip in pairs(site_exclude_ip)
|
||||
do
|
||||
cache:set(input_server_name .. "_exclude_ip_".._ip, true)
|
||||
for _, _ip in pairs(site_exclude_ip) do
|
||||
cache:set(input_server_name .. "_exclude_ip_" .. _ip, true)
|
||||
end
|
||||
end
|
||||
|
||||
-- set tag
|
||||
cache:set(load_key, true)
|
||||
return true
|
||||
end
|
||||
|
||||
--[[
|
||||
根据状态码判断是否排除
|
||||
检查当前请求的状态码是否在排除列表中
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function filter_status()
|
||||
if not auto_config['exclude_status'] then return false end
|
||||
local the_status = tostring(ngx.status)
|
||||
for _,v in ipairs(auto_config['exclude_status'])
|
||||
do
|
||||
for _, v in ipairs(auto_config['exclude_status']) do
|
||||
if the_status == v then
|
||||
return true
|
||||
end
|
||||
@@ -120,33 +149,51 @@ log_by_lua_block {
|
||||
return false
|
||||
end
|
||||
|
||||
--[[
|
||||
根据文件扩展名判断是否排除
|
||||
检查请求的 URI 是否以排除的扩展名结尾
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function exclude_extension()
|
||||
if not ngx.var.uri then return false end
|
||||
if not auto_config['exclude_extension'] then return false end
|
||||
for _,v in ipairs(auto_config['exclude_extension'])
|
||||
do
|
||||
if ngx.re.find(ngx.var.uri,"[.]"..v.."$",'ijo') then
|
||||
local uri = ngx.var.uri
|
||||
if not uri then return false end
|
||||
|
||||
local exclude_ext = auto_config['exclude_extension']
|
||||
if not exclude_ext then return false end
|
||||
|
||||
local lower_uri = string.lower(uri)
|
||||
for _, ext in ipairs(exclude_ext) do
|
||||
if string.find(lower_uri, "." .. ext .. "$", 1, true) then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
--[[
|
||||
根据 URL 判断是否排除
|
||||
支持精确匹配和正则匹配两种模式
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function exclude_url()
|
||||
if not ngx.var.uri then return false end
|
||||
if not ngx.var.request_uri then return false end
|
||||
if not auto_config['exclude_url'] then return false end
|
||||
local the_uri = string.sub(ngx.var.request_uri, 2)
|
||||
local url_conf = auto_config["exclude_url"]
|
||||
for i,conf in ipairs(url_conf)
|
||||
do
|
||||
local request_uri = ngx.var.request_uri
|
||||
if not request_uri then return false end
|
||||
|
||||
local url_conf = auto_config['exclude_url']
|
||||
if not url_conf then return false end
|
||||
|
||||
-- 去掉开头的 '/'
|
||||
local the_uri = string.sub(request_uri, 2)
|
||||
for _, conf in ipairs(url_conf) do
|
||||
local mode = conf["mode"]
|
||||
local url = conf["url"]
|
||||
if mode == "regular" then
|
||||
-- 正则匹配模式
|
||||
if ngx.re.find(the_uri, url, "ijo") then
|
||||
return true
|
||||
end
|
||||
else
|
||||
-- 精确匹配模式
|
||||
if the_uri == url then
|
||||
return true
|
||||
end
|
||||
@@ -155,513 +202,193 @@ log_by_lua_block {
|
||||
return false
|
||||
end
|
||||
|
||||
--[[
|
||||
根据 IP 判断是否排除
|
||||
先检查站点级排除 IP,再检查全局排除 IP
|
||||
@param input_server_name string 站点名称
|
||||
@param ip string 客户端 IP
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function exclude_ip(input_server_name, ip)
|
||||
-- 排除IP匹配,分网站单独的配置和全局配置两种方式
|
||||
local site_config = config[input_server_name]
|
||||
local server_exclude_ips = nil
|
||||
local check_server_exclude_ip = false
|
||||
if site_config then
|
||||
server_exclude_ips = site_config["exclude_ip"]
|
||||
if not server_exclude_ips then
|
||||
return false
|
||||
local server_exclude_ips = site_config["exclude_ip"]
|
||||
if server_exclude_ips then
|
||||
for _, _ip in pairs(server_exclude_ips) do
|
||||
if cache:get(input_server_name .. "_exclude_ip_" .. ip) then
|
||||
return true
|
||||
end
|
||||
for k, _ip in pairs(server_exclude_ips)
|
||||
do
|
||||
check_server_exclude_ip = true
|
||||
break
|
||||
end
|
||||
end
|
||||
-- D("server[" ..input_server_name.."]check exclude ip : "..tostring(check_server_exclude_ip))
|
||||
if check_server_exclude_ip then
|
||||
if cache:get(input_server_name .. "_exclude_ip_"..ip) then
|
||||
-- D("-Exclude server ip:"..ip)
|
||||
return true
|
||||
end
|
||||
else
|
||||
if cache:get("global_exclude_ip_"..ip) then
|
||||
-- D("*Excluded global ip:"..ip)
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
--------------------- exclude_func end ---------------------------
|
||||
|
||||
local function cache_logs_old(server_name)
|
||||
|
||||
-- make new id
|
||||
local new_id = C:get_last_id(server_name)
|
||||
|
||||
local excluded = false
|
||||
local ip = C:get_client_ip()
|
||||
excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(server_name, ip)
|
||||
|
||||
local ip_list = request_header["x-forwarded-for"]
|
||||
if ip and not ip_list then
|
||||
ip_list = ip
|
||||
end
|
||||
|
||||
local remote_addr = ngx.var.remote_addr
|
||||
if not string.find(ip_list, remote_addr) then
|
||||
if remote_addr then
|
||||
ip_list = ip_list .. "," .. remote_addr
|
||||
end
|
||||
return cache:get("global_exclude_ip_" .. ip) ~= nil
|
||||
end
|
||||
|
||||
-- local request_time = ngx.var.request_time
|
||||
local request_time = C:get_request_time()
|
||||
local client_port = ngx.var.remote_port
|
||||
local real_server_name = ngx.var.server_name
|
||||
local uri = ngx.var.uri
|
||||
local status_code = ngx.status
|
||||
local protocol = ngx.var.server_protocol
|
||||
local request_uri = ngx.var.request_uri
|
||||
local time_key = C:get_store_key()
|
||||
local method = method
|
||||
local body_length = C:get_length()
|
||||
local domain = C:get_domain()
|
||||
local referer = ngx.var.http_referer
|
||||
--[[
|
||||
日志采集函数
|
||||
采集请求数据,识别爬虫/客户端,计算统计指标,然后入队
|
||||
|
||||
local kv = {
|
||||
id=new_id,
|
||||
time_key=time_key,
|
||||
time=ngx.time(),
|
||||
ip=ip,
|
||||
domain=domain,
|
||||
server_name=server_name,
|
||||
real_server_name=real_server_name,
|
||||
method=method,
|
||||
status_code=status_code,
|
||||
uri=uri,
|
||||
request_uri=request_uri,
|
||||
body_length=body_length,
|
||||
referer=referer,
|
||||
user_agent=request_header['user-agent'],
|
||||
protocol=protocol,
|
||||
is_spider=0,
|
||||
request_time=request_time,
|
||||
excluded=excluded,
|
||||
request_headers='',
|
||||
ip_list=ip_list,
|
||||
client_port=client_port
|
||||
}
|
||||
|
||||
local request_stat_fields = "req=req+1,length=length+"..body_length
|
||||
local spider_stat_fields = "x"
|
||||
local client_stat_fields = "x"
|
||||
|
||||
if not excluded then
|
||||
|
||||
if status_code == 500 or (method=="POST" and config["record_post_args"] == true) or (status_code==403 and config["record_get_403_args"] == true) then
|
||||
local data = ""
|
||||
local ok, err = pcall(function() data = C:get_http_origin() end)
|
||||
if ok and not err then
|
||||
kv["request_headers"] = data
|
||||
end
|
||||
end
|
||||
|
||||
if ngx.re.find("500,501,502,503,504,505,506,507,509,510,400,401,402,403,404,405,406,407,408,409,410,411,412,413,414,415,416,417,418,421,422,423,424,425,426,449,451,499", tostring(status_code), "jo") then
|
||||
local field = "status_"..status_code
|
||||
request_stat_fields = request_stat_fields .. ","..field.."="..field.."+1"
|
||||
end
|
||||
|
||||
local lower_method = string.lower(method)
|
||||
if ngx.re.find("get,post,put,patch,delete", lower_method, "ijo") then
|
||||
local field = "http_"..lower_method
|
||||
request_stat_fields = request_stat_fields .. ","..field.."="..field.."+1"
|
||||
end
|
||||
|
||||
|
||||
local ipc = 0
|
||||
local pvc = 0
|
||||
local uvc = 0
|
||||
|
||||
local is_spider, request_spider, spider_index = C:match_spider(kv['user_agent'])
|
||||
if not is_spider then
|
||||
|
||||
client_stat_fields = C:match_client(kv['user_agent'])
|
||||
if not client_stat_fields or #client_stat_fields == 0 then
|
||||
client_stat_fields = request_stat_fields..",other=other+1"
|
||||
end
|
||||
|
||||
pvc, uvc = C:statistics_request(ip, is_spider,body_length)
|
||||
ipc = C:statistics_ipc(server_name,ip)
|
||||
else
|
||||
kv["is_spider"] = spider_index
|
||||
local field = "spider"
|
||||
spider_stat_fields = request_spider.."="..request_spider.."+"..1
|
||||
request_stat_fields = request_stat_fields .. ","..field.."="..field.."+"..1
|
||||
end
|
||||
|
||||
if ipc > 0 then
|
||||
request_stat_fields = request_stat_fields..",ip=ip+1"
|
||||
end
|
||||
if uvc > 0 then
|
||||
request_stat_fields = request_stat_fields..",uv=uv+1"
|
||||
end
|
||||
if pvc > 0 then
|
||||
request_stat_fields = request_stat_fields..",pv=pv+1"
|
||||
end
|
||||
end
|
||||
|
||||
local stat_fields = request_stat_fields..";"..client_stat_fields..";"..spider_stat_fields
|
||||
|
||||
cache_set(server_name, new_id, "stat_fields", stat_fields)
|
||||
cache_set(server_name, new_id, "log_kv", json.encode(kv))
|
||||
end
|
||||
执行流程:
|
||||
1. 获取日志 ID 和客户端 IP
|
||||
2. 判断是否需要排除(状态码/扩展名/URL/IP)
|
||||
3. 构建 IP 列表(含代理链)
|
||||
4. 采集请求数据(URI、状态码、响应大小等)
|
||||
5. 构建日志记录
|
||||
6. 识别爬虫/客户端类型
|
||||
7. 计算 PV/UV/IP 指标
|
||||
8. 将数据入队到共享内存
|
||||
|
||||
@param input_sn string 站点名称
|
||||
]]
|
||||
local function cache_logs(input_sn)
|
||||
|
||||
-- make new id
|
||||
-- 获取日志唯一 ID
|
||||
local new_id = C:get_last_id(input_sn)
|
||||
|
||||
local excluded = false
|
||||
-- 获取客户端真实 IP(支持 CDN 转发)
|
||||
local ip = C:get_client_ip()
|
||||
excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(input_sn, ip)
|
||||
-- 判断是否排除(状态码/扩展名/URL/IP)
|
||||
local excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(input_sn, ip)
|
||||
|
||||
-- 构建 IP 列表(包含代理链)
|
||||
local ip_list = request_header["x-forwarded-for"]
|
||||
if ip and not ip_list then
|
||||
ip_list = ip
|
||||
end
|
||||
|
||||
local remote_addr = ngx.var.remote_addr
|
||||
|
||||
-- 修复反向代理代过来的数据
|
||||
if "table" == type(ip_list) then
|
||||
ip_list = json.encode(ip_list)
|
||||
end
|
||||
|
||||
if not string.find(ip_list, remote_addr) then
|
||||
if remote_addr then
|
||||
if remote_addr and not string.find(ip_list, remote_addr, 1, true) then
|
||||
ip_list = ip_list .. "," .. remote_addr
|
||||
end
|
||||
end
|
||||
|
||||
-- local request_time = ngx.var.request_time
|
||||
local request_time = C:get_request_time()
|
||||
local client_port = ngx.var.remote_port
|
||||
local real_server_name = input_sn
|
||||
local uri = tostring(ngx.var.uri)
|
||||
local status_code = ngx.status
|
||||
local protocol = ngx.var.server_protocol
|
||||
local request_uri = ngx.var.request_uri
|
||||
local time_key = C:get_store_key()
|
||||
local method = method
|
||||
local body_length = C:get_length()
|
||||
local domain = C:get_domain()
|
||||
local referer = ngx.var.http_referer
|
||||
-- 采集请求数据
|
||||
local request_time = C:get_request_time() -- 请求耗时(毫秒)
|
||||
local client_port = ngx.var.remote_port -- 客户端端口
|
||||
local uri = tostring(ngx.var.uri) -- 请求路径
|
||||
local status_code = ngx.status -- 状态码
|
||||
local protocol = ngx.var.server_protocol -- 协议(HTTP/1.1 等)
|
||||
local request_uri = ngx.var.request_uri -- 完整请求路径(含参数)
|
||||
local body_length = C:get_length() -- 响应体大小
|
||||
local domain = C:get_domain() -- 域名
|
||||
local referer = ngx.var.http_referer -- 来源页
|
||||
local user_agent = request_header['user-agent'] -- 用户代理
|
||||
local now_time = ngx.time() -- 当前时间戳
|
||||
|
||||
-- 构建日志记录
|
||||
local kv = {
|
||||
id=new_id,
|
||||
time_key=time_key,
|
||||
time=ngx.time(),
|
||||
ip=ip,
|
||||
domain=domain,
|
||||
server_name=input_sn,
|
||||
real_server_name=real_server_name,
|
||||
method=method,
|
||||
status_code=status_code,
|
||||
uri=uri,
|
||||
request_uri=request_uri,
|
||||
body_length=body_length,
|
||||
referer=referer,
|
||||
user_agent=request_header['user-agent'],
|
||||
protocol=protocol,
|
||||
is_spider=0,
|
||||
request_time=request_time,
|
||||
excluded=excluded,
|
||||
request_headers='',
|
||||
ip_list=ip_list,
|
||||
client_port=client_port
|
||||
id = new_id,
|
||||
time_key = os.date("%Y%m%d%H", now_time), -- 小时级存储键
|
||||
time = now_time,
|
||||
ip = ip,
|
||||
domain = domain,
|
||||
server_name = input_sn,
|
||||
real_server_name = input_sn,
|
||||
method = method,
|
||||
status_code = status_code,
|
||||
uri = uri,
|
||||
request_uri = request_uri,
|
||||
body_length = body_length,
|
||||
referer = referer,
|
||||
user_agent = user_agent,
|
||||
protocol = protocol,
|
||||
is_spider = 0, -- 是否爬虫
|
||||
request_time = request_time,
|
||||
excluded = excluded, -- 是否排除
|
||||
request_headers = '', -- 请求头(异常时记录)
|
||||
ip_list = ip_list, -- IP 列表(含代理链)
|
||||
client_port = client_port
|
||||
}
|
||||
|
||||
-- C:D(json.encode(kv))
|
||||
local request_stat_fields = {req=1,length=body_length}
|
||||
|
||||
-- 初始化统计字段
|
||||
local request_stat_fields = {req = 1, length = body_length}
|
||||
local spider_stat_fields = {}
|
||||
local client_stat_fields = {}
|
||||
|
||||
-- 非排除请求的额外处理
|
||||
if not excluded then
|
||||
if status_code == 500 or (method=="POST" and config['global']["record_post_args"] == true) or (status_code==403 and config['global']["record_get_403_args"] == true) then
|
||||
local data = ""
|
||||
local ok, err = pcall(function() data = C:get_http_origin() end)
|
||||
if ok and not err then
|
||||
-- 记录异常请求的原始数据(500/POST/403)
|
||||
if status_code == 500 or (method == "POST" and config['global']["record_post_args"] == true) or (status_code == 403 and config['global']["record_get_403_args"] == true) then
|
||||
local ok, data = pcall(function() return C:get_http_origin() end)
|
||||
if ok and data then
|
||||
kv["request_headers"] = data
|
||||
else
|
||||
C:D("debug request_headers error:"..tostring(err))
|
||||
end
|
||||
end
|
||||
|
||||
if ngx.re.find("500,501,502,503,504,505,506,507,509,510,400,401,402,403,404,405,406,407,408,409,410,411,412,413,414,415,416,417,418,421,422,423,424,425,426,449,451,499", tostring(status_code), "jo") then
|
||||
local field = "status_"..status_code
|
||||
request_stat_fields[field] = 1
|
||||
-- 统计状态码
|
||||
if status_codes_to_log[tostring(status_code)] then
|
||||
request_stat_fields["status_" .. status_code] = 1
|
||||
end
|
||||
|
||||
-- D("method:"..method)
|
||||
-- 统计 HTTP 方法
|
||||
local lower_method = string.lower(method)
|
||||
if ngx.re.find("get,post,put,patch,delete", lower_method, "ijo") then
|
||||
local field = "http_"..lower_method
|
||||
request_stat_fields[field] = 1
|
||||
if http_methods[lower_method] then
|
||||
request_stat_fields["http_" .. lower_method] = 1
|
||||
end
|
||||
|
||||
|
||||
local ipc = 0
|
||||
local pvc = 0
|
||||
local uvc = 0
|
||||
|
||||
local is_spider, request_spider, spider_index = C:match_spider(kv['user_agent'])
|
||||
-- 识别爬虫/客户端
|
||||
local is_spider, request_spider, spider_index = C:match_spider(user_agent)
|
||||
if not is_spider then
|
||||
-- 非爬虫:识别客户端类型,计算 PV/UV/IP
|
||||
client_stat_fields = C:match_client_arr(user_agent)
|
||||
local pvc, uvc = C:statistics_request(ip, is_spider, body_length)
|
||||
local ipc = C:statistics_ipc(input_sn, ip)
|
||||
|
||||
client_stat_fields = C:match_client_arr(kv['user_agent'])
|
||||
pvc, uvc = C:statistics_request(ip, is_spider,body_length)
|
||||
ipc = C:statistics_ipc(input_sn,ip)
|
||||
if ipc > 0 then request_stat_fields["ip"] = 1 end
|
||||
if uvc > 0 then request_stat_fields["uv"] = 1 end
|
||||
if pvc > 0 then request_stat_fields["pv"] = 1 end
|
||||
else
|
||||
-- 爬虫:记录爬虫类型
|
||||
kv["is_spider"] = spider_index
|
||||
local field = "spider"
|
||||
spider_stat_fields[request_spider] = 1
|
||||
request_stat_fields[field] = 1
|
||||
end
|
||||
|
||||
if ipc > 0 then
|
||||
request_stat_fields["ip"] = 1
|
||||
end
|
||||
if uvc > 0 then
|
||||
request_stat_fields["uv"] = 1
|
||||
end
|
||||
if pvc > 0 then
|
||||
request_stat_fields["pv"] = 1
|
||||
request_stat_fields["spider"] = 1
|
||||
end
|
||||
end
|
||||
|
||||
local stat_fields = {
|
||||
request_stat_fields=request_stat_fields,
|
||||
client_stat_fields=client_stat_fields,
|
||||
spider_stat_fields=spider_stat_fields,
|
||||
}
|
||||
|
||||
-- 构建最终数据结构
|
||||
local data = {
|
||||
server_name=input_sn,
|
||||
stat_fields=stat_fields,
|
||||
log_kv=kv,
|
||||
server_name = input_sn,
|
||||
stat_fields = {
|
||||
request_stat_fields = request_stat_fields,
|
||||
client_stat_fields = client_stat_fields,
|
||||
spider_stat_fields = spider_stat_fields,
|
||||
},
|
||||
log_kv = kv,
|
||||
}
|
||||
|
||||
local push_data = json.encode(data)
|
||||
-- C:D(json.encode(push_data))
|
||||
local key = C:getTotalKey()
|
||||
ngx.shared.mw_total:rpush(key, push_data)
|
||||
end
|
||||
|
||||
local function store_logs_line(db, stmt, input_server_name, lineno)
|
||||
local logvalue = cache_get(input_server_name, lineno, "log_kv")
|
||||
if not logvalue then return false end
|
||||
local logline = json.decode(logvalue)
|
||||
|
||||
local time = logline["time"]
|
||||
local id = logline["id"]
|
||||
local protocol = logline["protocol"]
|
||||
local client_port = logline["client_port"]
|
||||
local status_code = logline["status_code"]
|
||||
local uri = logline["uri"]
|
||||
local request_uri = logline["request_uri"]
|
||||
local method = logline["method"]
|
||||
local body_length = logline["body_length"]
|
||||
local referer = logline["referer"]
|
||||
local ip = logline["ip"]
|
||||
local ip_list = logline["ip_list"]
|
||||
local request_time = logline["request_time"]
|
||||
local is_spider = logline["is_spider"]
|
||||
local domain = logline["domain"]
|
||||
local server_name = logline["server_name"]
|
||||
local user_agent = logline["user_agent"]
|
||||
local request_headers = logline["request_headers"]
|
||||
local excluded = logline["excluded"]
|
||||
|
||||
local request_stat_fields = nil
|
||||
local client_stat_fields = nil
|
||||
local spider_stat_fields = nil
|
||||
local stat_fields = cache_get(input_server_name, id, "stat_fields")
|
||||
if stat_fields == nil then
|
||||
-- D("Log stat fields is nil.")
|
||||
-- D("Logdata:"..logvalue)
|
||||
else
|
||||
stat_fields = C:split(stat_fields, ";")
|
||||
request_stat_fields = stat_fields[1]
|
||||
client_stat_fields = stat_fields[2]
|
||||
spider_stat_fields = stat_fields[3]
|
||||
|
||||
if "x" == client_stat_fields then
|
||||
client_stat_fields = nil
|
||||
end
|
||||
|
||||
if "x" == spider_stat_fields then
|
||||
spider_stat_fields = nil
|
||||
end
|
||||
end
|
||||
|
||||
local time_key = logline["time_key"]
|
||||
if not excluded then
|
||||
stmt:bind_names{
|
||||
time=time,
|
||||
ip=ip,
|
||||
domain=domain,
|
||||
server_name=server_name,
|
||||
method=method,
|
||||
status_code=status_code,
|
||||
uri=request_uri,
|
||||
body_length=body_length,
|
||||
referer=referer,
|
||||
user_agent=user_agent,
|
||||
protocol=protocol,
|
||||
request_time=request_time,
|
||||
is_spider=is_spider,
|
||||
request_headers=request_headers,
|
||||
ip_list=ip_list,
|
||||
client_port=client_port,
|
||||
}
|
||||
|
||||
local res, err = stmt:step()
|
||||
if tostring(res) == "5" then
|
||||
-- D("step res:"..tostring(res))
|
||||
-- D("step err:"..tostring(err))
|
||||
-- D("the step database connection is busy, so it will be stored later.")
|
||||
return false
|
||||
end
|
||||
stmt:reset()
|
||||
-- D("store_logs_line ok")
|
||||
|
||||
C:update_stat( db, "client_stat", time_key, client_stat_fields)
|
||||
C:update_stat( db, "spider_stat", time_key, spider_stat_fields)
|
||||
-- C:D("stat ok")
|
||||
|
||||
-- only count non spider requests
|
||||
local ok, err = C:statistics_uri(db, request_uri, ngx.md5(request_uri), body_length)
|
||||
local ok, err = C:statistics_ip(db, ip, body_length)
|
||||
end
|
||||
|
||||
C:update_stat( db, "request_stat", time_key, request_stat_fields)
|
||||
return true
|
||||
end
|
||||
|
||||
local function store_logs(input_sn)
|
||||
if C:is_migrating(input_sn) == true then
|
||||
-- D("migrating...")
|
||||
return
|
||||
end
|
||||
|
||||
local last_insert_id_key = input_sn.."_last_id"
|
||||
local store_start_id_key = input_sn.."_store_start"
|
||||
local last_id = cache:get(last_insert_id_key)
|
||||
local store_start = cache:get(store_start_id_key)
|
||||
if store_start == nil then
|
||||
store_start = 1
|
||||
end
|
||||
local store_end = last_id
|
||||
if store_end == nil then
|
||||
store_end = 1
|
||||
end
|
||||
|
||||
local worker_id = ngx.worker.id()
|
||||
if C:is_working(input_sn) then
|
||||
-- D("other workers are being stored, please store later.")
|
||||
-- cache:delete(flush_data_key)
|
||||
return true
|
||||
end
|
||||
C:lock_working(input_sn)
|
||||
|
||||
local log_dir = "{$SERVER_APP}/logs"
|
||||
local db_path = log_dir .. '/' .. input_sn .. "/logs.db"
|
||||
local db, err = sqlite3.open(db_path)
|
||||
|
||||
if tostring(err) ~= 'nil' then
|
||||
C:D("sqlite3 open error:"..tostring(err))
|
||||
return true
|
||||
end
|
||||
|
||||
db:exec([[PRAGMA synchronous = 0]])
|
||||
db:exec([[PRAGMA page_size = 4096]])
|
||||
db:exec([[PRAGMA journal_mode = wal]])
|
||||
db:exec([[PRAGMA journal_size_limit = 1073741824]])
|
||||
|
||||
local stmt2 = nil
|
||||
if db ~= nil then
|
||||
stmt2 = 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 db == nil or stmt2 == nil then
|
||||
-- D("web data db error")
|
||||
-- cache:set(storing_key, false)
|
||||
if db and db:isopen() then
|
||||
db:close()
|
||||
end
|
||||
return true
|
||||
end
|
||||
|
||||
status, errorString = db:exec([[BEGIN TRANSACTION]])
|
||||
|
||||
C:clean_stats(db, input_sn)
|
||||
|
||||
if store_end >= store_start then
|
||||
for i=store_start, store_end, 1 do
|
||||
-- D("store_start:"..store_start..":store_end:".. store_end)
|
||||
if store_logs_line(db, stmt2, input_sn, i) then
|
||||
cache_clear(input_sn, i, "log_kv")
|
||||
cache_clear(input_sn, i, "stat_fields")
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
local res, err = stmt2:finalize()
|
||||
if tostring(res) == "5" then
|
||||
C:D("Finalize res:"..tostring(res)..",Finalize err:"..tostring(err))
|
||||
end
|
||||
|
||||
local now_date = os.date("*t")
|
||||
local save_day = config['global']["save_day"]
|
||||
local save_date_timestamp = os.time{year=now_date.year,
|
||||
month=now_date.month, day=now_date.day-save_day, hour=0}
|
||||
-- delete expire data
|
||||
db:exec("DELETE FROM web_logs WHERE time<"..tostring(save_date_timestamp))
|
||||
|
||||
local res, err = db:execute([[COMMIT]])
|
||||
if db and db:isopen() then
|
||||
db:close()
|
||||
end
|
||||
|
||||
cache:set(store_start_id_key, store_end+1)
|
||||
C:unlock_working(input_sn)
|
||||
-- 入队到共享内存队列
|
||||
cache:rpush(total_key, json.encode(data))
|
||||
end
|
||||
|
||||
--[[
|
||||
应用入口函数
|
||||
加载排除 IP 列表,然后采集日志
|
||||
]]
|
||||
local function run_app()
|
||||
-- C:D("------------ debug start ------------")
|
||||
init_var()
|
||||
|
||||
load_global_exclude_ip()
|
||||
load_exclude_ip(server_name)
|
||||
|
||||
cache_logs(server_name)
|
||||
-- C:cron()
|
||||
|
||||
-- cache_logs_old(server_name)
|
||||
-- store_logs(server_name)
|
||||
-- C:D("------------ debug end -------------")
|
||||
end
|
||||
|
||||
|
||||
--[[
|
||||
应用入口包装函数(带错误捕获)
|
||||
在调试模式下,使用 pcall 捕获异常并记录日志
|
||||
]]
|
||||
local function run_app_ok()
|
||||
if not debug_mode then return run_app() end
|
||||
|
||||
local presult, err = pcall( function() run_app() end)
|
||||
local presult, err = pcall(function() run_app() end)
|
||||
if not presult then
|
||||
C:D("debug error on :"..tostring(err))
|
||||
C:D("debug error on :" .. tostring(err))
|
||||
return true
|
||||
end
|
||||
end
|
||||
|
||||
-- 执行应用
|
||||
return run_app_ok()
|
||||
}
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ import io
|
||||
import os
|
||||
import time
|
||||
import json
|
||||
import fcntl
|
||||
|
||||
web_dir = os.getcwd() + "/web"
|
||||
if os.path.exists(web_dir):
|
||||
@@ -44,7 +45,6 @@ def getGlobalConf():
|
||||
|
||||
|
||||
def pSqliteDb(dbname='web_logs', site_name='unset', fn="logs"):
|
||||
|
||||
db_dir = getServerDir() + '/logs/' + site_name
|
||||
if not os.path.exists(db_dir):
|
||||
mw.execShell('mkdir -p ' + db_dir)
|
||||
@@ -67,143 +67,196 @@ def pSqliteDb(dbname='web_logs', site_name='unset', fn="logs"):
|
||||
|
||||
conn.autoTextFactory()
|
||||
|
||||
# conn.text_factory = lambda x: str(x, encoding="utf-8", errors='ignore')
|
||||
# conn.text_factory = lambda x: unicode(x, "utf-8", "ignore")
|
||||
return conn
|
||||
|
||||
|
||||
def migrateSiteHotLogs(site_name, query_date):
|
||||
print(site_name, query_date)
|
||||
start_time = time.time()
|
||||
print(f"[{site_name}] 开始迁移日志...")
|
||||
|
||||
migrating_flag = getServerDir() + "/logs/%s/migrating" % site_name
|
||||
hot_db = getServerDir() + "/logs/%s/logs.db" % site_name
|
||||
hot_db_tmp = getServerDir() + "/logs/%s/logs_tmp.db" % site_name
|
||||
history_logs_db = getServerDir() + "/logs/%s/history_logs.db" % site_name
|
||||
# 1. copy to tmp file
|
||||
|
||||
if not os.path.exists(hot_db):
|
||||
print(f"[{site_name}] 热日志数据库不存在,跳过")
|
||||
return mw.returnMsg(True, f"{site_name} logs not exist, skip")
|
||||
|
||||
# 检查是否正在迁移
|
||||
if os.path.exists(migrating_flag):
|
||||
print(f"[{site_name}] 正在迁移中,跳过")
|
||||
return mw.returnMsg(True, f"{site_name} is migrating, skip")
|
||||
|
||||
# 1. 备份到临时文件(使用文件直接拷贝)
|
||||
try:
|
||||
import shutil
|
||||
print("coping {} to {} ...".format(hot_db, hot_db_tmp))
|
||||
print(f"[{site_name}] 备份 {hot_db} -> {hot_db_tmp} ...")
|
||||
mw.writeFile(migrating_flag, "yes")
|
||||
time.sleep(3)
|
||||
time.sleep(0.5)
|
||||
shutil.copy(hot_db, hot_db_tmp)
|
||||
if not os.path.exists(hot_db_tmp):
|
||||
return mw.returnMsg(False, "migrating fail, copy tmp file!")
|
||||
except:
|
||||
return mw.returnMsg(False, "{} migrating fail.".format(site_name))
|
||||
finally:
|
||||
return mw.returnMsg(False, f"{site_name} migrating fail, copy tmp file!")
|
||||
except Exception as e:
|
||||
if os.path.exists(migrating_flag):
|
||||
os.remove(migrating_flag)
|
||||
return mw.returnMsg(False, f"{site_name} migrating fail: {e}")
|
||||
|
||||
# 2. 从临时备份中迁移热日志数据到历史日志
|
||||
print("begin tmp to hot log data ...")
|
||||
# 2. 从临时备份中迁移热日志数据到历史日志(批量插入)
|
||||
try:
|
||||
print("history file: {}".format(history_logs_db))
|
||||
logs_conn = pSqliteDb('web_log', site_name, 'logs_tmp')
|
||||
history_logs_conn = pSqliteDb('web_log', site_name, 'history_logs')
|
||||
|
||||
hot_db_columns = logs_conn.originExecute(
|
||||
"PRAGMA table_info([web_logs])")
|
||||
_columns = ",".join([c[1] for c in hot_db_columns if c[1] != "id"])
|
||||
query_start = 0
|
||||
_columns = [c[1] for c in hot_db_columns if c[1] != "id"]
|
||||
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 = "select {} from web_logs where time<{}".format(
|
||||
_columns, todayUt)
|
||||
logs_sql = f"select {columns_str} from web_logs where time<{todayUt}"
|
||||
selector = logs_conn.originExecute(logs_sql)
|
||||
log = selector.fetchone()
|
||||
while log:
|
||||
params = ""
|
||||
for field in log:
|
||||
if params:
|
||||
params += ","
|
||||
if field is None:
|
||||
field = "\'\'"
|
||||
elif type(field) == str:
|
||||
field = "\'" + field.replace("\'", "\"") + "\'"
|
||||
params += str(field)
|
||||
insert_sql = "insert into web_logs(" + \
|
||||
_columns + ") values(" + params + ")"
|
||||
history_logs_conn.execute(insert_sql)
|
||||
log = selector.fetchone()
|
||||
|
||||
print("sorting historical data, this action takes a long time...")
|
||||
history_logs_conn.execute("VACUUM;")
|
||||
# 批量插入配置
|
||||
batch_size = 10000
|
||||
batch_count = 0
|
||||
total_count = 0
|
||||
insert_sql = f"insert into web_logs({columns_str}) values({placeholders})"
|
||||
|
||||
while True:
|
||||
logs = selector.fetchmany(batch_size)
|
||||
if not logs:
|
||||
break
|
||||
|
||||
batch_count += 1
|
||||
total_count += len(logs)
|
||||
# 使用原生 executemany 批量插入
|
||||
history_logs_conn.executemany(insert_sql, logs)
|
||||
history_logs_conn.commit()
|
||||
|
||||
if batch_count % 10 == 0:
|
||||
elapsed = time.time() - start_time
|
||||
print(f"[{site_name}] 已迁移 {total_count} 条记录, 耗时 {elapsed:.2f}s")
|
||||
|
||||
print(f"[{site_name}] 共迁移 {total_count} 条记录")
|
||||
|
||||
# 清理历史数据
|
||||
gcfg = getGlobalConf()
|
||||
save_day = gcfg['global']["save_day"]
|
||||
print("delete historical data {} days ago...".format(save_day))
|
||||
print(f"[{site_name}] 删除 {save_day} 天前的历史数据...")
|
||||
|
||||
time_now = time.localtime()
|
||||
save_timestamp = time.mktime((time_now.tm_year, time_now.tm_mon, time_now.tm_mday - save_day, 0, 0, 0, 0, 0, 0))
|
||||
delete_sql = "delete from web_logs where time <= {}".format(
|
||||
save_timestamp)
|
||||
print('delete history_logs')
|
||||
print(delete_sql)
|
||||
delete_sql = f"delete from web_logs where time <= {save_timestamp}"
|
||||
history_logs_conn.execute(delete_sql)
|
||||
history_logs_conn.commit()
|
||||
|
||||
# 3. delete merged data and clean up statistics
|
||||
print("delete merged thermal data...")
|
||||
history_logs_conn.execute("VACUUM;")
|
||||
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}")
|
||||
return mw.returnMsg(False, f"{site_name} logs migrate error: {e}")
|
||||
|
||||
# 3. 删除已迁移的数据并清理统计(批量删除)
|
||||
try:
|
||||
if os.path.exists(migrating_flag):
|
||||
os.remove(migrating_flag)
|
||||
|
||||
mw.writeFile(migrating_flag, "yes")
|
||||
|
||||
hot_db_conn = pSqliteDb('web_logs', site_name)
|
||||
del_hot_log = "delete from web_logs where time<{}".format(todayUt)
|
||||
print(del_hot_log)
|
||||
r = hot_db_conn.execute(del_hot_log)
|
||||
print("delete:", r)
|
||||
print("deleting statistics over 180 days...")
|
||||
|
||||
# 分批删除热日志
|
||||
del_hot_log = f"delete from web_logs where time<{todayUt}"
|
||||
print(f"[{site_name}] 删除已迁移的热日志...")
|
||||
hot_db_conn.execute(del_hot_log)
|
||||
|
||||
# 删除过期统计数据
|
||||
print(f"[{site_name}] 删除180天前的统计数据...")
|
||||
save_time_key = time.strftime(
|
||||
'%Y%m%d00', time.localtime(time.time() - 180 * 86400))
|
||||
|
||||
del_request_stat_sql = "delete from request_stat where time<={}".format(
|
||||
save_time_key)
|
||||
del_request_stat_sql = f"delete from request_stat where time<={save_time_key}"
|
||||
hot_db_conn.execute(del_request_stat_sql)
|
||||
hot_db_conn.execute(f"delete from spider_stat where time<={save_time_key}")
|
||||
hot_db_conn.execute(f"delete from client_stat where time<={save_time_key}")
|
||||
hot_db_conn.execute(f"delete from referer_stat where time<={save_time_key}")
|
||||
|
||||
hot_db_conn.execute(
|
||||
"delete from spider_stat where time<={}".format(save_time_key))
|
||||
hot_db_conn.execute(
|
||||
"delete from client_stat where time<={}".format(save_time_key))
|
||||
hot_db_conn.execute(
|
||||
"delete from referer_stat where time<={}".format(save_time_key))
|
||||
hot_db_conn.commit()
|
||||
print("clean up the hot database...")
|
||||
print(f"[{site_name}] 压缩热数据库...")
|
||||
hot_db_conn.execute("VACUUM;")
|
||||
hot_db_conn.commit()
|
||||
|
||||
except Exception as e:
|
||||
print(f"[{site_name}] delete hot logs error: {e}")
|
||||
finally:
|
||||
if os.path.exists(migrating_flag):
|
||||
os.remove(migrating_flag)
|
||||
except Exception as e:
|
||||
if site_name:
|
||||
print("{} logs to history error:{}".format(site_name, e))
|
||||
else:
|
||||
print("logs to history error:{}".format(e))
|
||||
finally:
|
||||
if os.path.exists(hot_db_tmp):
|
||||
os.remove(hot_db_tmp)
|
||||
|
||||
print("{} logs migrate ok.".format(site_name))
|
||||
elapsed = time.time() - start_time
|
||||
print(f"[{site_name}] 日志迁移完成,耗时 {elapsed:.2f}s")
|
||||
|
||||
if not mw.isAppleSystem():
|
||||
mw.execShell("chown -R www:www " + getServerDir())
|
||||
|
||||
mw.opWeb('restart')
|
||||
return mw.returnMsg(True, "{} logs migrate ok".format(site_name))
|
||||
return mw.returnMsg(True, f"{site_name} logs migrate ok")
|
||||
|
||||
|
||||
def migrateHotLogs(query_date="today"):
|
||||
print("begin migrate hot logs")
|
||||
# 使用文件锁防止并发迁移
|
||||
lock_file = getServerDir() + "/logs/migrate_hot_logs.lock"
|
||||
lock_fd = None
|
||||
|
||||
try:
|
||||
lock_fd = open(lock_file, 'w')
|
||||
fcntl.flock(lock_fd.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
|
||||
print("开始迁移热日志...")
|
||||
sites = mw.M('sites').field('name').order("add_time").select()
|
||||
|
||||
unset_site = {"name": "unset"}
|
||||
sites.append(unset_site)
|
||||
|
||||
# migrateSiteHotLogs('t1.cn', query_date)
|
||||
total_sites = len(sites)
|
||||
success_count = 0
|
||||
fail_count = 0
|
||||
|
||||
for site_info in sites:
|
||||
# print(site_info['name'])
|
||||
for i, site_info in enumerate(sites):
|
||||
site_name = site_info["name"]
|
||||
print(f"\n[{i+1}/{total_sites}] 处理站点: {site_name}")
|
||||
migrate_res = migrateSiteHotLogs(site_name, query_date)
|
||||
if not migrate_res["status"]:
|
||||
print(migrate_res["msg"])
|
||||
print("end migrate hot logs")
|
||||
if migrate_res["status"]:
|
||||
success_count += 1
|
||||
else:
|
||||
fail_count += 1
|
||||
print(f"[{site_name}] 迁移失败: {migrate_res['msg']}")
|
||||
|
||||
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:
|
||||
print("正在迁移中,跳过")
|
||||
return mw.returnMsg(True, "正在迁移中,跳过")
|
||||
except Exception as e:
|
||||
print(f"migrateHotLogs error: {e}")
|
||||
return mw.returnMsg(False, f"migrateHotLogs error: {e}")
|
||||
finally:
|
||||
if lock_fd:
|
||||
try:
|
||||
fcntl.flock(lock_fd.fileno(), fcntl.LOCK_UN)
|
||||
lock_fd.close()
|
||||
if os.path.exists(lock_file):
|
||||
os.remove(lock_file)
|
||||
except:
|
||||
pass
|
||||
+5
-1
@@ -335,9 +335,13 @@ class Sql():
|
||||
except Exception as ex:
|
||||
return "error: " + str(ex)
|
||||
|
||||
def commit(self):
|
||||
def executemany(self, sql, args):
|
||||
self.__DB_CONN.executemany(sql, args)
|
||||
self.__close()
|
||||
|
||||
def commit(self):
|
||||
self.__DB_CONN.commit()
|
||||
self.__close()
|
||||
|
||||
def save(self, keys, param):
|
||||
# 更新数据
|
||||
|
||||
@@ -201,6 +201,10 @@ def getCommonFile():
|
||||
}
|
||||
return data
|
||||
|
||||
def getSysTTL():
|
||||
f = "/proc/sys/net/ipv4/ip_default_ttl"
|
||||
return readFile(f).strip()
|
||||
|
||||
def checkCert(certPath='ssl/certificate.pem'):
|
||||
# 验证证书
|
||||
openssl = '/usr/bin/openssl'
|
||||
|
||||
Reference in New Issue
Block a user