forked from huawei/openGauss-server
fix (dbmind): fix some bugs on the index-advisor and add extract_log.py
This commit is contained in:
parent
9ffe07c287
commit
3f9b68d3b5
|
|
@ -0,0 +1,83 @@
|
|||
import re
|
||||
import os
|
||||
import argparse
|
||||
|
||||
|
||||
SQL_TYPE = ['select', 'delete', 'insert', 'update']
|
||||
|
||||
|
||||
def output_valid_sql(sql):
|
||||
if any(tp in sql.lower() for tp in SQL_TYPE[1:]):
|
||||
return sql if sql.endswith('; ') else sql + ';'
|
||||
elif SQL_TYPE[0] in sql.lower() and 'from ' in sql.lower():
|
||||
return sql if sql.endswith('; ') else sql + ';'
|
||||
return ''
|
||||
|
||||
|
||||
def extract_sql_from_log(log_path):
|
||||
files = os.listdir(log_path)
|
||||
for file in files:
|
||||
if not os.path.isdir(file) and '.swap' not in file:
|
||||
with open(log_path + "/" + file, mode='r') as f:
|
||||
line = f.readline()
|
||||
sql = ''
|
||||
statement_flag = False
|
||||
execute_flag = False
|
||||
|
||||
while line:
|
||||
# Identify statement scene
|
||||
if re.search('statement: ', line, re.IGNORECASE):
|
||||
statement_flag = True
|
||||
if output_valid_sql(sql):
|
||||
yield output_valid_sql(sql)
|
||||
sql = re.search(r'statement: (.*)', line.strip(),
|
||||
re.IGNORECASE).group(1) + ' '
|
||||
line = f.readline()
|
||||
|
||||
# Identify execute statement scene
|
||||
elif re.search(r'execute .*:', line, re.IGNORECASE):
|
||||
if output_valid_sql(sql):
|
||||
yield output_valid_sql(sql)
|
||||
execute_flag = True
|
||||
sql = re.search(r'execute .*: (.*)', line.strip(), re.IGNORECASE).group(1)
|
||||
line = f.readline()
|
||||
else:
|
||||
if statement_flag:
|
||||
if re.match(r'^\t', line):
|
||||
sql += line.strip('\t\n')
|
||||
else:
|
||||
statement_flag = False
|
||||
if output_valid_sql(sql):
|
||||
yield output_valid_sql(sql)
|
||||
sql = ''
|
||||
if execute_flag and re.search(r'parameters: ', line, re.IGNORECASE):
|
||||
execute_flag = False
|
||||
param_list = re.search(r'parameters: (.*)', line.strip(),
|
||||
re.IGNORECASE).group(1).split(', ')
|
||||
param_list = list(param.split('=', 1) for param in param_list)
|
||||
param_list.sort(key=lambda x: int(x[0].strip(' $')),
|
||||
reverse=True)
|
||||
for item in param_list:
|
||||
if len(item[1].strip()) >= 256:
|
||||
sql = sql.replace(item[0].strip(), "''")
|
||||
else:
|
||||
sql = sql.replace(item[0].strip(), item[1].strip())
|
||||
yield output_valid_sql(sql)
|
||||
sql = ''
|
||||
line = f.readline()
|
||||
|
||||
|
||||
def main():
|
||||
arg_parser = argparse.ArgumentParser()
|
||||
arg_parser.add_argument("l", help="The path of the log file that needs to be parsed.")
|
||||
arg_parser.add_argument("f", help="The output path of the extracted file.")
|
||||
args = arg_parser.parse_args()
|
||||
sqls = extract_sql_from_log(args.l)
|
||||
with open(args.f, 'w') as file:
|
||||
for sql in sqls:
|
||||
file.write(sql + '\n')
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
|
||||
|
|
@ -104,12 +104,12 @@ def run_shell_cmd(target_sql_list):
|
|||
|
||||
def run_shell_sql_cmd(sql_file):
|
||||
cmd = BASE_CMD + ' -f ./' + sql_file
|
||||
try:
|
||||
ret = subprocess.check_output(shlex.split(cmd), stderr=subprocess.STDOUT)
|
||||
except subprocess.CalledProcessError as e:
|
||||
print(e.output, file=sys.stderr)
|
||||
|
||||
proc = subprocess.Popen(shlex.split(cmd), stdout=subprocess.PIPE, stderr=subprocess.PIPE)
|
||||
(stdout, stderr) = proc.communicate()
|
||||
stdout, stderr = stdout.decode(), stderr.decode()
|
||||
|
||||
return stdout
|
||||
return ret.decode()
|
||||
|
||||
|
||||
def green(text):
|
||||
|
|
@ -219,18 +219,19 @@ def estimate_workload_cost_file(workload, index_config=None):
|
|||
found_plan = False
|
||||
hypo_index = False
|
||||
with open(sql_file, 'w') as file:
|
||||
if SCHEMA:
|
||||
file.write('SET current_schema = %s;\n' % SCHEMA)
|
||||
if index_config:
|
||||
# create hypo-indexes
|
||||
if SCHEMA:
|
||||
file.write('SET current_schema = %s;\n' % SCHEMA)
|
||||
file.write('SET enable_hypo_index = on;\n')
|
||||
for index in index_config:
|
||||
file.write("SELECT hypopg_create_index('CREATE INDEX ON %s(%s)');\n" %
|
||||
(index.table, index.columns))
|
||||
if ENABLE_MULTI_NODE:
|
||||
file.write('set enable_fast_query_shipping = off;\n')
|
||||
file.write("set explain_perf_mode = 'normal'; \n")
|
||||
for query in workload:
|
||||
file.write('set explain_perf_mode = 'normal'; EXPLAIN ' + query.statement + ';\n')
|
||||
file.write('EXPLAIN ' + query.statement + ';\n')
|
||||
|
||||
result = run_shell_sql_cmd(sql_file).split('\n')
|
||||
if os.path.exists(sql_file):
|
||||
|
|
@ -328,7 +329,7 @@ def query_index_check(query, query_index_dict):
|
|||
if columns != '':
|
||||
sql_list.append("SELECT hypopg_create_index('CREATE INDEX ON %s(%s)')" %
|
||||
(table, columns))
|
||||
sql_list.append('explain ' + query)
|
||||
sql_list.append("set explain_perf_mode = 'normal'; explain " + query)
|
||||
sql_list.append('SELECT hypopg_reset_index()')
|
||||
result = run_shell_cmd(sql_list).split('\n')
|
||||
|
||||
|
|
@ -442,7 +443,7 @@ def generate_candidate_indexes_file(workload):
|
|||
for column in columns:
|
||||
if column == "":
|
||||
continue
|
||||
if column not in index_dict[table]:
|
||||
if column not in index_dict[table] and not re.match(r'\s*,\s*$', column):
|
||||
print("table: ", table, "columns: ", column)
|
||||
index_dict[table].add(column)
|
||||
candidate_indexes.append(IndexItem(table, column))
|
||||
|
|
@ -570,7 +571,7 @@ def infer_workload_cost(workload, config, atomic_config_total):
|
|||
return total_cost
|
||||
|
||||
|
||||
def simple_index_advisor(input_path):
|
||||
def simple_index_advisor(input_path, max_index_num):
|
||||
workload = workload_compression(input_path)
|
||||
print_header_boundary(" Generate candidate indexes ")
|
||||
candidate_indexes = generate_candidate_indexes_file(workload)
|
||||
|
|
@ -587,6 +588,7 @@ def simple_index_advisor(input_path):
|
|||
|
||||
# filter out duplicate index
|
||||
final_index_set = {}
|
||||
final_index_list = []
|
||||
for index in candidate_indexes:
|
||||
picked = True
|
||||
cols = set(index.columns.split(','))
|
||||
|
|
@ -600,23 +602,20 @@ def simple_index_advisor(input_path):
|
|||
if len(cols.difference(pre_cols)) == 0:
|
||||
picked = False
|
||||
break
|
||||
if picked:
|
||||
if picked and index.benefit > 0:
|
||||
final_index_set[index.table].append(index)
|
||||
|
||||
[final_index_list.extend(item) for item in list(final_index_set.values())]
|
||||
final_index_list = sorted(final_index_list, key=lambda item: item.benefit, reverse=True)
|
||||
cnt = 0
|
||||
index_current_storage = 0
|
||||
for table, indexs in final_index_set.items():
|
||||
for index in indexs:
|
||||
if cnt == MAX_INDEX_NUM:
|
||||
break
|
||||
if MAX_INDEX_STORAGE and (index_current_storage + index.storage) > MAX_INDEX_STORAGE:
|
||||
continue
|
||||
index_current_storage += index.storage
|
||||
print("create index ind" + str(cnt) + " on " + table + "(" + index.columns + ");")
|
||||
cnt += 1
|
||||
else:
|
||||
for index in final_index_list:
|
||||
if max_index_num and cnt == max_index_num:
|
||||
break
|
||||
if MAX_INDEX_STORAGE and (index_current_storage + index.storage) > MAX_INDEX_STORAGE:
|
||||
continue
|
||||
break
|
||||
index_current_storage += index.storage
|
||||
print("create index ind" + str(cnt) + " on " + index.table + "(" + index.columns + ");")
|
||||
cnt += 1
|
||||
|
||||
|
||||
def greedy_determine_opt_config(workload, atomic_config_total, candidate_indexes):
|
||||
|
|
@ -728,7 +727,7 @@ def main():
|
|||
if args.multi_iter_mode:
|
||||
complex_index_advisor(args.f)
|
||||
else:
|
||||
simple_index_advisor(args.f)
|
||||
simple_index_advisor(args.f, args.max_index_num)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
|
|
|
|||
Loading…
Reference in New Issue