This commit is contained in:
2023-02-15 22:04:30 +08:00
34 changed files with 1268 additions and 686 deletions
+6
View File
@@ -59,6 +59,10 @@ data/benchmark
!nyctx100.csv !nyctx100.csv
!network.csv !network.csv
!test_complex.csv !test_complex.csv
data/electricity*
data/covtype*
data/phishing*
data/power*
*.out *.out
*.asm *.asm
!mmw.so !mmw.so
@@ -83,3 +87,5 @@ udf*.hpp
*.ipynb *.ipynb
saved_procedures/** saved_procedures/**
procedures/** procedures/**
.mypy_cache
__pycache__
+3
View File
@@ -0,0 +1,3 @@
[submodule "paper"]
path = paper
url = https://github.com/sunyinqi0508/AQueryPaper
+1 -1
View File
@@ -10,7 +10,7 @@ RUN export OS_VER=`cat /etc/os-release | grep VERSION_CODENAME` &&\
RUN wget --output-document=/etc/apt/trusted.gpg.d/monetdb.gpg https://dev.monetdb.org/downloads/MonetDB-GPG-KEY.gpg RUN wget --output-document=/etc/apt/trusted.gpg.d/monetdb.gpg https://dev.monetdb.org/downloads/MonetDB-GPG-KEY.gpg
RUN apt update && apt install -y python3 python3-pip clang-14 libmonetdbe-dev git RUN apt update && apt install -y python3 python3-pip clang-14 libmonetdbe-dev libmonetdb-client-dev monetdb5-sql-dev git
RUN git clone https://github.com/sunyinqi0508/AQuery2 RUN git clone https://github.com/sunyinqi0508/AQuery2
+5 -3
View File
@@ -2,6 +2,7 @@ OS_SUPPORT =
MonetDB_LIB = MonetDB_LIB =
MonetDB_INC = MonetDB_INC =
Defines = Defines =
CC = $(CXX) -xc
CXXFLAGS = --std=c++2a CXXFLAGS = --std=c++2a
ifeq ($(AQ_DEBUG), 1) ifeq ($(AQ_DEBUG), 1)
OPTFLAGS = -g3 #-fsanitize=address OPTFLAGS = -g3 #-fsanitize=address
@@ -17,7 +18,7 @@ COMPILER = $(strip $(_COMPILER))
LIBTOOL = ar rcs LIBTOOL = ar rcs
USELIB_FLAG = -Wl,--whole-archive,libaquery.a -Wl,-no-whole-archive USELIB_FLAG = -Wl,--whole-archive,libaquery.a -Wl,-no-whole-archive
LIBAQ_SRC = server/monetdb_conn.cpp server/libaquery.cpp LIBAQ_SRC = server/monetdb_conn.cpp server/libaquery.cpp
LIBAQ_OBJ = monetdb_conn.o libaquery.o LIBAQ_OBJ = monetdb_conn.o libaquery.o monetdb_ext.o
SEMANTIC_INTERPOSITION = -fno-semantic-interposition SEMANTIC_INTERPOSITION = -fno-semantic-interposition
RANLIB = ranlib RANLIB = ranlib
_LINKER_BINARY = $(shell `$(CXX) -print-prog-name=ld` -v 2>&1 | grep -q LLVM && echo lld || echo ld) _LINKER_BINARY = $(shell `$(CXX) -print-prog-name=ld` -v 2>&1 | grep -q LLVM && echo lld || echo ld)
@@ -43,7 +44,7 @@ else
LIBTOOL = gcc-ar rcs LIBTOOL = gcc-ar rcs
endif endif
endif endif
OPTFLAGS += $(SEMANTIC_INTERPOSITION) LINKFLAGS += $(SEMANTIC_INTERPOSITION)
ifeq ($(PCH), 1) ifeq ($(PCH), 1)
PCHFLAGS = -include server/pch.hpp PCHFLAGS = -include server/pch.hpp
@@ -82,7 +83,7 @@ else
MonetDB_INC += $(AQ_MONETDB_INC) MonetDB_INC += $(AQ_MONETDB_INC)
MonetDB_INC += -I/usr/local/include/monetdb -I/usr/include/monetdb MonetDB_INC += -I/usr/local/include/monetdb -I/usr/include/monetdb
endif endif
MonetDB_LIB += -lmonetdbe MonetDB_LIB += -lmonetdbe -lmonetdbsql -lbat
endif endif
ifeq ($(THREADING),1) ifeq ($(THREADING),1)
@@ -128,6 +129,7 @@ pch:
$(CXX) -x c++-header server/pch.hpp $(FPIC) $(CXXFLAGS) $(CXX) -x c++-header server/pch.hpp $(FPIC) $(CXXFLAGS)
libaquery: libaquery:
$(CXX) -c $(FPIC) $(PCHFLAGS) $(LIBAQ_SRC) $(OS_SUPPORT) $(CXXFLAGS) &&\ $(CXX) -c $(FPIC) $(PCHFLAGS) $(LIBAQ_SRC) $(OS_SUPPORT) $(CXXFLAGS) &&\
$(CC) -c server/monetdb_ext.c $(OPTFLAGS) $(MonetDB_INC) &&\
$(LIBTOOL) libaquery.a $(LIBAQ_OBJ) &&\ $(LIBTOOL) libaquery.a $(LIBAQ_OBJ) &&\
$(RANLIB) libaquery.a $(RANLIB) libaquery.a
+1 -1
View File
@@ -2,7 +2,7 @@
## GLOBAL CONFIGURATION FLAGS ## GLOBAL CONFIGURATION FLAGS
version_string = '0.6.0a' version_string = '0.7.0a'
add_path_to_ldpath = True add_path_to_ldpath = True
rebuild_backend = False rebuild_backend = False
run_backend = True run_backend = True
+5 -41
View File
@@ -5,20 +5,18 @@
# You can obtain one at http://mozilla.org/MPL/2.0/. # You can obtain one at http://mozilla.org/MPL/2.0/.
# #
# Contact: Kyle Lahnakoski (kyle@lahnakoski.com) # Contact: Kyle Lahnakoski (kyle@lahnakoski.com)
# # Bill Sun 2022 - 2023
from __future__ import absolute_import, division, unicode_literals from __future__ import absolute_import, division, unicode_literals
import json import json
from threading import Lock from threading import Lock
from aquery_parser.sql_parser import scrub from aquery_parser.parser import scrub
from aquery_parser.utils import ansi_string, simple_op, normal_op from aquery_parser.utils import ansi_string, simple_op, normal_op
import aquery_parser.parser
parse_locker = Lock() # ENSURE ONLY ONE PARSING AT A TIME parse_locker = Lock() # ENSURE ONLY ONE PARSING AT A TIME
common_parser = None common_parser = None
mysql_parser = None
sqlserver_parser = None
SQL_NULL = {"null": {}} SQL_NULL = {"null": {}}
@@ -33,44 +31,10 @@ def parse(sql, null=SQL_NULL, calls=simple_op):
with parse_locker: with parse_locker:
if not common_parser: if not common_parser:
common_parser = sql_parser.common_parser() common_parser = aquery_parser.parser.common_parser()
result = _parse(common_parser, sql, null, calls) result = _parse(common_parser, sql, null, calls)
return result return result
def parse_mysql(sql, null=SQL_NULL, calls=simple_op):
"""
PARSE MySQL ASSUME DOUBLE QUOTED STRINGS ARE LITERALS
:param sql: String of SQL
:param null: What value to use as NULL (default is the null function `{"null":{}}`)
:return: parse tree
"""
global mysql_parser
with parse_locker:
if not mysql_parser:
mysql_parser = sql_parser.mysql_parser()
return _parse(mysql_parser, sql, null, calls)
def parse_sqlserver(sql, null=SQL_NULL, calls=simple_op):
"""
PARSE MySQL ASSUME DOUBLE QUOTED STRINGS ARE LITERALS
:param sql: String of SQL
:param null: What value to use as NULL (default is the null function `{"null":{}}`)
:return: parse tree
"""
global sqlserver_parser
with parse_locker:
if not sqlserver_parser:
sqlserver_parser = sql_parser.sqlserver_parser()
return _parse(sqlserver_parser, sql, null, calls)
parse_bigquery = parse_mysql
def _parse(parser, sql, null, calls): def _parse(parser, sql, null, calls):
utils.null_locations = [] utils.null_locations = []
utils.scrub_op = calls utils.scrub_op = calls
@@ -85,4 +49,4 @@ def _parse(parser, sql, null, calls):
_ = json.dumps _ = json.dumps
__all__ = ["parse", "format", "parse_mysql", "parse_bigquery", "normal_op", "simple_op"] __all__ = ["parse", "format", "normal_op", "simple_op"]
+1 -1
View File
@@ -5,7 +5,7 @@
# You can obtain one at http://mozilla.org/MPL/2.0/. # You can obtain one at http://mozilla.org/MPL/2.0/.
# #
# Contact: Kyle Lahnakoski (kyle@lahnakoski.com) # Contact: Kyle Lahnakoski (kyle@lahnakoski.com)
# # Bill Sun 2022 - 2023
# SQL CONSTANTS # SQL CONSTANTS
from mo_parsing import * from mo_parsing import *
@@ -5,7 +5,7 @@
# You can obtain one at http://mozilla.org/MPL/2.0/. # You can obtain one at http://mozilla.org/MPL/2.0/.
# #
# Contact: Kyle Lahnakoski (kyle@lahnakoski.com) # Contact: Kyle Lahnakoski (kyle@lahnakoski.com)
# # Bill Sun 2022 - 2023
from sre_parse import WHITESPACE from sre_parse import WHITESPACE
@@ -28,37 +28,12 @@ simple_ident = Regex(simple_ident.__regex__()[1])
def common_parser(): def common_parser():
combined_ident = Combine(delimited_list( combined_ident = Combine(delimited_list(
ansi_ident | mysql_backtick_ident | simple_ident, separator=".", combine=True, ansi_ident | aquery_backtick_ident | simple_ident, separator=".", combine=True,
)).set_parser_name("identifier") )).set_parser_name("identifier")
return parser(ansi_string | mysql_doublequote_string, combined_ident) return parser(ansi_string | aquery_doublequote_string, combined_ident)
def parser(literal_string, ident):
def mysql_parser():
mysql_string = ansi_string | mysql_doublequote_string
mysql_ident = Combine(delimited_list(
mysql_backtick_ident | sqlserver_ident | simple_ident,
separator=".",
combine=True,
)).set_parser_name("mysql identifier")
return parser(mysql_string, mysql_ident)
def sqlserver_parser():
combined_ident = Combine(delimited_list(
ansi_ident
| mysql_backtick_ident
| sqlserver_ident
| Word(FIRST_IDENT_CHAR, IDENT_CHAR),
separator=".",
combine=True,
)).set_parser_name("identifier")
return parser(ansi_string, combined_ident, sqlserver=True)
def parser(literal_string, ident, sqlserver=False):
with Whitespace() as engine: with Whitespace() as engine:
engine.add_ignore(Literal("--") + restOfLine) engine.add_ignore(Literal("--") + restOfLine)
engine.add_ignore(Literal("#") + restOfLine) engine.add_ignore(Literal("#") + restOfLine)
@@ -184,8 +159,6 @@ def parser(literal_string, ident, sqlserver=False):
) )
) )
if not sqlserver:
# SQL SERVER DOES NOT SUPPORT [] FOR ARRAY CONSTRUCTION (USED FOR IDENTIFIERS)
create_array = ( create_array = (
Literal("[") + delimited_list(Group(expr))("args") + Literal("]") Literal("[") + delimited_list(Group(expr))("args") + Literal("]")
| create_array | create_array
+1 -1
View File
@@ -5,7 +5,7 @@
# You can obtain one at http://mozilla.org/MPL/2.0/. # You can obtain one at http://mozilla.org/MPL/2.0/.
# #
# Contact: Kyle Lahnakoski (kyle@lahnakoski.com) # Contact: Kyle Lahnakoski (kyle@lahnakoski.com)
# # Bill Sun 2022 - 2023
# KNOWN TYPES # KNOWN TYPES
+3 -4
View File
@@ -5,7 +5,7 @@
# You can obtain one at http://mozilla.org/MPL/2.0/. # You can obtain one at http://mozilla.org/MPL/2.0/.
# #
# Contact: Kyle Lahnakoski (kyle@lahnakoski.com) # Contact: Kyle Lahnakoski (kyle@lahnakoski.com)
# # Bill Sun 2022 - 2023
import ast import ast
@@ -610,9 +610,8 @@ hex_num = (
# STRINGS # STRINGS
ansi_string = Regex(r"\'(\'\'|[^'])*\'") / to_string ansi_string = Regex(r"\'(\'\'|[^'])*\'") / to_string
mysql_doublequote_string = Regex(r'\"(\"\"|[^"])*\"') / to_string aquery_doublequote_string = Regex(r'\"(\"\"|[^"])*\"') / to_string
# BASIC IDENTIFIERS # BASIC IDENTIFIERS
ansi_ident = Regex(r'\"(\"\"|[^"])*\"') / unquote ansi_ident = Regex(r'\"(\"\"|[^"])*\"') / unquote
mysql_backtick_ident = Regex(r"\`(\`\`|[^`])*\`") / unquote aquery_backtick_ident = Regex(r"\`(\`\`|[^`])*\`") / unquote
sqlserver_ident = Regex(r"\[(\]\]|[^\]])*\]") / unquote
+2
View File
@@ -0,0 +1,2 @@
make snippet_uselib
cp ./dll.so procedures/q70.so
Submodule
+1
Submodule paper added at fa4e3f5a06
+15 -7
View File
@@ -5,14 +5,14 @@ from typing import List
name : str = input('Filename (in path ./procedures/<filename>.aqp):') name : str = input('Filename (in path ./procedures/<filename>.aqp):')
def write(): def write():
s : str = input() s : str = input('Enter queries: empty line to stop. \n')
qs : List[str] = [] qs : List[str] = []
while(len(s) and not s.startswith('S')): while(len(s) and not s.startswith('S')):
qs.append(s) qs.append(s)
s = input() s = input()
ms : int = int(input()) ms : int = int(input('number of modules:'))
with open(f'./procedures/{name}.aqp', 'wb') as fp: with open(f'./procedures/{name}.aqp', 'wb') as fp:
fp.write(struct.pack("I", len(qs) + (ms > 0))) fp.write(struct.pack("I", len(qs) + (ms > 0)))
@@ -27,21 +27,29 @@ def write():
fp.write(b'\x00') fp.write(b'\x00')
def read(): def read(cmd : str):
rc = len(cmd) > 1 and cmd[1] == 'c'
clip = ''
with open(f'./procedures/{name}.aqp', 'rb') as fp: with open(f'./procedures/{name}.aqp', 'rb') as fp:
nq = struct.unpack("I", fp.read(4))[0] nq = struct.unpack("I", fp.read(4))[0]
ms = struct.unpack("I", fp.read(4))[0] ms = struct.unpack("I", fp.read(4))[0]
qs = fp.read().split(b'\x00') qs = fp.read().split(b'\x00')
print(f'Procedure {name}, {nq} queries, {ms} modules:') print(f'Procedure {name}, {nq} queries, {ms} modules:')
for q in qs: for q in qs:
print(' ' + q.decode('utf-8')) q = q.decode('utf-8').strip()
if q:
q = f'"{q}",' if rc else f'\t{q}'
print(q)
clip += q + '\n'
if rc and not input('copy to clipboard?').lower().startswith('n'):
import pyperclip
pyperclip.copy(clip)
if __name__ == '__main__': if __name__ == '__main__':
while True: while True:
cmd = input("r for read, w for write: ") cmd = input("r for read, rc to read c_str, w for write: ")
if cmd.lower().startswith('r'): if cmd.lower().startswith('r'):
read() read(cmd.lower())
break break
elif cmd.lower().startswith('w'): elif cmd.lower().startswith('w'):
write() write()
+44
View File
@@ -0,0 +1,44 @@
import os
sep = os.sep
# toggles
dataset = 'power' # [covtype, electricity, mixed, phishing, power]
use_threadpool = False # True
# environments
input_prefix = f'data{sep}{dataset}_orig'
output_prefix = f'data{sep}{dataset}'
sep_field = b','
sep_subfield = b';'
lst_files = os.listdir(input_prefix)
# lst_files.sort()
try:
os.mkdir(output_prefix)
except FileExistsError:
pass
def process(f : str):
filename = input_prefix + sep + f
ofilename = output_prefix + sep + f[:-3] + 'csv'
with open(filename, 'rb') as ifile:
icontents = ifile.read()
with open(ofilename, 'wb') as ofile:
ofile.write(b'\n')
for l in icontents.splitlines():
fields = l.strip().split(b' ')
subfields = fields[:-1]
ol = ( # fields[0] + sep_field +
sep_subfield.join(subfields) +
sep_field + fields[-1] + b'\n')
ofile.write(ol)
if not use_threadpool:
for f in lst_files:
process(f)
elif __name__ == '__main__':
from multiprocessing import Pool
with Pool(8) as tp:
tp.map(process, lst_files)
+12 -9
View File
@@ -9,26 +9,29 @@ struct DR;
struct DT; struct DT;
//enum Evaluation {gini, entropy, logLoss}; class DecisionTree
{
class DecisionTree{
public: public:
DT *DTree = nullptr; DT *DTree = nullptr;
int maxHeight; double minIG;
long maxHeight;
long feature; long feature;
long maxFeature; long maxFeature;
long seed; bool isRF;
long classes; long classes;
int *Sparse; int *Sparse;
double forgetRate; double forgetRate;
double increaseRate;
double initialIR;
Evaluation evalue; Evaluation evalue;
long Rebuild; long Rebuild;
long roundNo; long roundNo;
long called; long called;
long retain; long retain;
long lastT;
long lastAll;
DecisionTree(int hight, long f, int* sparse, double forget, long maxFeature, long noClasses, Evaluation e, long r, long rb); DecisionTree(long f, int *sparse, double forget, long maxFeature, long noClasses, Evaluation e);
void Stablelize(); void Stablelize();
@@ -38,9 +41,9 @@ minEval findMinGiniDense(double** data, long* result, long* totalT, long size, l
minEval findMinGiniSparse(double **data, long *result, long *totalT, long size, long col, DT *current); minEval findMinGiniSparse(double **data, long *result, long *totalT, long size, long col, DT *current);
minEval incrementalMinGiniDense(double** data, long* result, long size, long col, long*** count, double** record, long* max, long newCount, long forgetSize, bool isRoot); minEval incrementalMinGiniDense(double **data, long *result, long size, long col, long ***count, double **record, long *max, long newCount, long forgetSize, double **forgottenData, long *forgottenClass);
minEval incrementalMinGiniSparse(double** dataNew, long* resultNew, long sizeNew, long sizeOld, DT* current, long col, long forgetSize, bool isRoot); minEval incrementalMinGiniSparse(double **dataNew, long *resultNew, long sizeNew, long sizeOld, DT *current, long col, long forgetSize, double **forgottenData, long *forgottenClass);
long *fitThenPredict(double **trainData, long *trainResult, long trainSize, double **testData, long testSize); long *fitThenPredict(double **trainData, long *trainResult, long trainSize, double **testData, long testSize);
+14 -70
View File
@@ -26,7 +26,7 @@ minEval giniSparse(double** data, long* result, long* d, long size, long col, lo
double gini1, gini2; double gini1, gini2;
double c; double c;
long l, r; long l, r;
for(i=0; i<size; i++){ for(i=0; i<size-1; i++){
c = data[d[i]][col]; c = data[d[i]][col];
if(c==max)break; if(c==max)break;
count[result[d[i]]]++; count[result[d[i]]]++;
@@ -62,7 +62,7 @@ minEval entropySparse(double** data, long* result, long* d, long size, long col,
double entropy1, entropy2; double entropy1, entropy2;
double c; double c;
long l, r; long l, r;
for(i=0; i<size; i++){ for(i=0; i<size-1; i++){
c = data[d[i]][col]; c = data[d[i]][col];
if(c==max)break; if(c==max)break;
count[result[d[i]]]++; count[result[d[i]]]++;
@@ -73,8 +73,8 @@ minEval entropySparse(double** data, long* result, long* d, long size, long col,
for(j=0;j<classes;j++){ for(j=0;j<classes;j++){
l = count[j]; l = count[j];
r = totalT[j]-l; r = totalT[j]-l;
entropy1 -= ((double)l/total)*log((double)l/total); if(l!=0)entropy1 -= ((double)l/total)*log((double)l/total);
entropy2 -= ((double)r/(size-total))*log((double)r/(size-total)); if(r!=0)entropy2 -= ((double)r/(size-total))*log((double)r/(size-total));
} }
entropy1 = entropy1*total/size + entropy2*(size-total)/size; entropy1 = entropy1*total/size + entropy2*(size-total)/size;
if(ret.eval>entropy1){ if(ret.eval>entropy1){
@@ -140,8 +140,8 @@ minEval entropySparseIncremental(long sizeTotal, long classes, double* newSorted
for(j=0;j<classes;j++){ for(j=0;j<classes;j++){
l = count[j]; l = count[j];
r = T[j]-l; r = T[j]-l;
e1 -= ((double)l/total)*log((double)l/total); if(l!=0)e1 -= ((double)l/total)*log((double)l/total);
e2 -= ((double)r/(sizeTotal-total))*log((double)r/(sizeTotal-total)); if(r!=0)e2 -= ((double)r/(sizeTotal-total))*log((double)r/(sizeTotal-total));
} }
e1 = e1*total/sizeTotal + e2*(sizeTotal-total)/sizeTotal; e1 = e1*total/sizeTotal + e2*(sizeTotal-total)/sizeTotal;
if(ret.eval>e1){ if(ret.eval>e1){
@@ -159,9 +159,9 @@ minEval giniDense(long max, long size, long classes, long** rem, long* d, double
double gini1, gini2; double gini1, gini2;
long *t, *t2, *r, *r2, i, j; long *t, *t2, *r, *r2, i, j;
for(i=0;i<max;i++){ for(i=0;i<max;i++){
t = rem[d[i]]; t = rem[i];
if(i>0){ if(i>0){
t2 = rem[d[i-1]]; t2 = rem[i-1];
for(j=0;j<=classes;j++){ for(j=0;j<=classes;j++){
t[j]+=t2[j]; t[j]+=t2[j];
} }
@@ -179,7 +179,7 @@ minEval giniDense(long max, long size, long classes, long** rem, long* d, double
gini1 = (gini1*t[classes])/size + (gini2*(size-t[classes]))/size; gini1 = (gini1*t[classes])/size + (gini2*(size-t[classes]))/size;
if(gini1<ret.eval){ if(gini1<ret.eval){
ret.eval = gini1; ret.eval = gini1;
ret.value = record[d[i]]; ret.value = record[i];
ret.left = t[classes]; ret.left = t[classes];
} }
} }
@@ -193,9 +193,9 @@ minEval entropyDense(long max, long size, long classes, long** rem, long* d, dou
double entropy1, entropy2; double entropy1, entropy2;
long *t, *t2, *r, *r2, i, j; long *t, *t2, *r, *r2, i, j;
for(i=0;i<max;i++){ for(i=0;i<max;i++){
t = rem[d[i]]; t = rem[i];
if(i>0){ if(i>0){
t2 = rem[d[i-1]]; t2 = rem[i-1];
for(j=0;j<=classes;j++){ for(j=0;j<=classes;j++){
t[j]+=t2[j]; t[j]+=t2[j];
} }
@@ -207,71 +207,15 @@ minEval entropyDense(long max, long size, long classes, long** rem, long* d, dou
long l, r; long l, r;
l = t[j]; l = t[j];
r = totalT[j]-l; r = totalT[j]-l;
entropy1 -= ((double)l/t[classes])*log((double)l/t[classes]); if(l!=0)entropy1 -= ((double)l/t[classes])*log((double)l/t[classes]);
entropy2 -= ((double)r/(size-t[classes]))*log((double)r/(size-t[classes])); if(r!=0)entropy2 -= ((double)r/(size-t[classes]))*log((double)r/(size-t[classes]));
} }
entropy1 = entropy1*t[classes]/size + entropy2*(size-t[classes])/size; entropy1 = entropy1*t[classes]/size + entropy2*(size-t[classes])/size;
if(entropy1<ret.eval){ if(entropy1<ret.eval){
ret.eval = entropy1; ret.eval = entropy1;
ret.value = record[d[i]]; ret.value = record[i];
ret.left = t[classes]; ret.left = t[classes];
} }
} }
return ret; return ret;
} }
minEval giniDenseIncremental(long max, double* record, long** count, long classes, long newSize, long* T){
double gini1, gini2;
minEval ret;
long i, j;
ret.eval = DBL_MAX;
for(i=0; i<max; i++){
if(count[i][classes]==newSize){
continue;
}
gini1 = 1.0;
gini2 = 1.0;
for(j=0;j<classes;j++){
long l, r;
l = count[i][j];
r = T[j]-l;
gini1 -= pow((double)l/count[i][classes], 2);
gini2 -= pow((double)r/(newSize-count[i][classes]), 2);
}
gini1 = gini1*count[i][classes]/newSize + gini2*((newSize-count[i][classes]))/newSize;
if(gini1<ret.eval){
ret.eval = gini1;
ret.value = record[i];
}
}
return ret;
}
minEval entropyDenseIncremental(long max, double* record, long** count, long classes, long newSize, long* T){
double entropy1, entropy2;
minEval ret;
long i, j;
ret.eval = DBL_MAX;
for(i=0; i<max; i++){
if(count[i][classes]==newSize or count[i][classes]==0){
continue;
}
entropy1 = 0;
entropy2 = 0;
for(j=0;j<classes;j++){
long l, r;
l = count[i][j];
r = T[j]-l;
entropy1 -= ((double)l/count[i][classes])*log((double)l/count[i][classes]);
entropy2 -= (double)r/(newSize-count[i][classes])*log((double)r/(newSize-count[i][classes]));
}
entropy1 = entropy1*count[i][classes]/newSize + entropy2*((newSize-count[i][classes]))/newSize;
if(entropy1<ret.eval){
ret.eval = entropy1;
ret.value = record[i];
}
}
return ret;
}
-4
View File
@@ -17,8 +17,4 @@ minEval giniDense(long max, long size, long classes, long** rem, long* d, double
minEval entropyDense(long max, long size, long classes, long** rem, long* d, double* record, long* totalT); minEval entropyDense(long max, long size, long classes, long** rem, long* d, double* record, long* totalT);
minEval giniDenseIncremental(long max, double* record, long** count, long classes, long newSize, long* T);
minEval entropyDenseIncremental(long max, double* record, long** count, long classes, long newSize, long* T);
#endif #endif
+182 -54
View File
@@ -2,7 +2,27 @@
#include <stdlib.h> #include <stdlib.h>
#include <stdio.h> #include <stdio.h>
#include <ctime> #include <ctime>
#include <math.h>
#include <algorithm>
#include <boost/math/distributions/students_t.hpp>
#include <random>
long poisson(int Lambda)
{
int k = 0;
long double p = 1.0;
long double l = exp(-Lambda);
srand((long)clock());
while(p>=l)
{
double u = (double)(rand()%10000)/10000;
p *= u;
k++;
}
if (k>11)k=11;
return k-1;
}
struct DT{ struct DT{
int height; int height;
long* featureId; long* featureId;
@@ -30,62 +50,126 @@ struct DT{
long size = 0;// Size of the dataset long size = 0;// Size of the dataset
}; };
RandomForest::RandomForest(long mTree, long actTree, long rTime, int h, long feature, int* s, double forg, long maxF, long noC, Evaluation eval, long r, long rb){ RandomForest::RandomForest(long mTree, long feature, int* s, double forg, long noC, Evaluation eval, bool b, double t){
srand((long)clock()); srand((long)clock());
Rebuild = rb; bagging = b;
if(actTree<1)actTree=1; activeTree = mTree;
noTree = actTree;
activeTree = actTree;
treePointer = 0;
if(mTree<actTree)mTree=activeTree;
maxTree = mTree; maxTree = mTree;
if(rTime<=0)rTime=1; allT = new long[mTree];
rotateTime = rTime;
timer = 0;
retain = r;
tThresh=t;
lastT = -2;
lastAll = 0;
long i; long i;
height = h;
f = feature; f = feature;
sparse = new int[f]; sparse = new int[f];
for(i=0; i<f; i++)sparse[i]=s[i]; for(i=0; i<f; i++)sparse[i]=s[i];
forget = forg; forget = forg;
maxFeature = maxF;
noClasses = noC; noClasses = noC;
e = eval; e = eval;
minF = floor(sqrt((double)f))+2;
if(minF>f)minF=f;
DTrees = (DecisionTree**)malloc(mTree*sizeof(DecisionTree*)); DTrees = (DecisionTree**)malloc(mTree*sizeof(DecisionTree*));
for(i=0; i<mTree; i++){ for(i=0; i<maxTree; i++){
if(i<actTree){ DTrees[i] = new DecisionTree(f, sparse, forget, minF+rand()%(f+1-minF), noClasses, e);
DTrees[i] = new DecisionTree(height, f, sparse, forget, maxFeature, noClasses, e, r, rb); DTrees[i]->isRF=true;
}
else{
DTrees[i]=nullptr;
}
} }
} }
void RandomForest::fit(double** data, long* result, long size){ void RandomForest::fit(double** data, long* result, long size){
if(timer==rotateTime and maxTree!=activeTree){ long i, j, k, l;
Rotate();
timer=0;
}
long i, j, k;
double** newData; double** newData;
long* newResult; long* newResult;
for(i=0; i<activeTree; i++){ long localT = 0;
newData = new double*[size]; int stale = 0;
newResult = new long[size]; if(lastT==-2){
lastT=-1;
}else{
for(i=0; i<maxTree; i++)allT[i] = 0;
for(i=0; i<size; i++){
if(Test(data[i], result[i])==result[i])localT++;
}
long localAll = size;
if(lastT>=0){
double lastSm = (double)lastT/lastAll;
double localSm = (double)localT/localAll;
double lastSd = sqrt(pow((1.0-lastSm),2)*lastT+pow(lastSm,2)*(lastAll-lastT)/(lastAll-1));
double localSd = sqrt(pow((1.0-localSm),2)*localT+pow(localSm,2)*(localAll-localT)/(localAll-1));
double v = lastAll+localAll-2;
double sp = sqrt(((lastAll-1) * lastSd * lastSd + (localAll-1) * localSd * localSd) / v);
double q;
double t = lastSm-localSm;
if(sp==0){q = 1;}
else{
t = t/(sp*sqrt(1.0/lastAll+1.0/localAll));
boost::math::students_t dist(v);
double c = cdf(dist, t);
q = cdf(complement(dist, fabs(t)));
}
if(q<=tThresh){
lastT += localT;
lastAll += localAll;
}else if(t<0){
lastT = localT;
lastAll = localAll;
}else{
double newAcc = (double)localT/localAll;
double lastAcc= (double)lastT/lastAll;
stale = floor((newAcc-lastAcc)/(lastAcc)*maxTree);
lastT = localT;
lastAll = localAll;
}
}else{
lastT = localT;
lastAll = localAll;
}
}
Rotate(stale);
for(i=0; i<maxTree; i++){
long times;
if(bagging)times = poisson(6);
else times=1;
if(times==0)continue;
newData = (double**)malloc(sizeof(double*)*size*times);
newResult = (long*)malloc(sizeof(long)*size*times);
long c = 0;
for(j = 0; j<size*times; j++){
long jj;
if(bagging) jj = rand()%size;
else jj=j;
newData[j] = (double*)malloc((f+1)*sizeof(double));
for(l=0; l<f; l++){
newData[j][l] = data[jj][l];
}
newData[j][f] = 0;
newResult[j] = result[jj];
}
DTrees[i]->fit(newData, newResult, size*times);
}
/*for(i=0; i<maxTree; i++){
//backupTrees[i]->retain = 10*size;
//if(backupTrees[i]==nullptr) continue;
long times;
//times = poisson(posMean);
//if(times==0)continue;
times=1;
newData = (double**)malloc(sizeof(double*)*size*times);
newResult = (long*)malloc(sizeof(long)*size*times);
long c = 0;
for(j = 0; j<size; j++){ for(j = 0; j<size; j++){
newData[j] = new double[f]; long jj = rand()%size;
for(k=0; k<f; k++){ jj=j;
newData[j][k] = data[j][k]; for(k=0; k<times; k++){
newData[j*times+k] = (double*)malloc((f+1)*sizeof(double));
for(l=0; l<f; l++){
newData[j*times+k][l] = data[jj][l];
} }
newResult[j] = result[j]; newData[j*times+k][f] = 0;
newResult[j*times+k] = result[jj];
} }
DTrees[(i+treePointer)%maxTree]->fit(newData, newResult, size);
} }
timer++; backupTrees[i]->fit(newData, newResult, size*times);
}*/
} }
long* RandomForest::fitThenPredict(double** trainData, long* trainResult, long trainSize, double** testData, long testSize){ long* RandomForest::fitThenPredict(double** trainData, long* trainResult, long trainSize, double** testData, long testSize){
@@ -97,36 +181,80 @@ long* RandomForest::fitThenPredict(double** trainData, long* trainResult, long t
return testResult; return testResult;
} }
void RandomForest::Rotate(){ void RandomForest::Rotate(long stale){
if(noTree==maxTree){ long i, j, k;
DTrees[(treePointer+activeTree)%maxTree]->Free(); long minIndex = -1;
delete DTrees[(treePointer+activeTree)%maxTree]; if(stale>=0)return;
}else{ else{
noTree++; stale = std::min(stale, maxTree);
stale*=-1;
while(stale>0){
long currentMin = 2147483647;
for(i = 0; i<maxTree; i++){
if(allT[i]<currentMin){
currentMin=allT[i];
minIndex=i;
}
}
stale--;
if(minIndex<0)break;
allT[minIndex] = 2147483647;
double** newData;
long* newResult;
long size = 0;
long lastT2 = 0;
size = DTrees[minIndex]->DTree->size;
newData = (double**)malloc(sizeof(double*)*size);
newResult = (long*)malloc(sizeof(long)*size);
for(j = 0; j<size; j++){
newData[j] = (double*)malloc(sizeof(double)*(f+1));
for(k=0; k<f; k++){
newData[j][k] = DTrees[minIndex]->DTree->dataRecord[j][k];
}
newData[j][f] = 0;
newResult[j] = DTrees[minIndex]->DTree->resultRecord[j];
}
DTrees[minIndex]->Stablelize();
DTrees[minIndex]->Free();
delete DTrees[minIndex];
DTrees[minIndex] = new DecisionTree(f, sparse, forget, minF+rand()%(f+1-minF), noClasses, e);
DTrees[minIndex]->isRF=true;
DTrees[minIndex]->fit(newData, newResult, size);
for(j=0; j<size; j++){
if(DTrees[minIndex]->Test(newData[j], DTrees[minIndex]->DTree)==newResult[j])lastT2++;
}
DTrees[minIndex]->lastAll=size;
DTrees[minIndex]->lastT=lastT2;
} }
DTrees[(treePointer+activeTree)%maxTree] = new DecisionTree(height, f, sparse, forget, maxFeature, noClasses, e, retain, Rebuild);
long size = DTrees[(treePointer+activeTree-1)%maxTree]->DTree->size;
double** newData = new double*[size];
long* newResult = new long[size];
for(long j = 0; j<size; j++){
newData[j] = new double[f];
for(long k=0; k<f; k++){
newData[j][k] = DTrees[(treePointer+activeTree-1)%maxTree]->DTree->dataRecord[j][k];
} }
newResult[j] = DTrees[(treePointer+activeTree-1)%maxTree]->DTree->resultRecord[j];
} }
DTrees[(treePointer+activeTree)%maxTree]->fit(newData, newResult, size);
DTrees[treePointer]->Stablelize(); long RandomForest::Test(double* data, long result){
if(++treePointer==maxTree)treePointer=0; long i;
long predict[noClasses];
for(i=0; i<noClasses; i++){
predict[i]=0;
}
for(i=0; i<maxTree; i++){
long tmp = DTrees[i]->Test(data, DTrees[i]->DTree);
predict[tmp]++;
if(tmp==result)allT[i]++;
} }
long ret = 0;
for(i=1; i<noClasses; i++){
if(predict[i]>predict[ret])ret = i;
}
return ret;
}
long RandomForest::Test(double* data){ long RandomForest::Test(double* data){
long i; long i;
long predict[noClasses]; long predict[noClasses];
for(i=0; i<noClasses; i++)predict[i]=0; for(i=0; i<noClasses; i++)predict[i]=0;
for(i=0; i<noTree; i++){ for(i=0; i<maxTree; i++){
predict[DTrees[i]->Test(data, DTrees[i]->DTree)]++; predict[DTrees[i]->Test(data, DTrees[i]->DTree)]++;
} }
+12 -11
View File
@@ -14,33 +14,34 @@ struct DT;
class RandomForest{ class RandomForest{
public: public:
long noTree;
long maxTree; long maxTree;
long activeTree; long activeTree;
long treePointer; long* allT;
long rotateTime; double tThresh;
long timer; DecisionTree** DTrees;
long retain; DecisionTree** backupTrees;
DecisionTree** DTrees = nullptr;
long height; long height;
long Rebuild; bool bagging;
long f; long f;
int* sparse; int* sparse;
double forget; double forget;
long maxFeature;
long noClasses; long noClasses;
Evaluation e; Evaluation e;
long lastT;
long lastAll;
int minF;
RandomForest(long maxTree, long f, int* sparse, double forget, long noClasses=2, Evaluation e=Evaluation::entropy, bool b=false, double tThresh=0.05);
RandomForest(long maxTree, long activeTree, long rotateTime, int height, long f, int* sparse, double forget, long maxFeature=0, long noClasses=2, Evaluation e=Evaluation::gini, long r=-1, long rb=2147483647);
void fit(double** data, long* result, long size); void fit(double** data, long* result, long size);
long* fitThenPredict(double** trainData, long* trainResult, long trainSize, double** testData, long testSize); long* fitThenPredict(double** trainData, long* trainResult, long trainSize, double** testData, long testSize);
void Rotate(); void Rotate(long stale);
long Test(double* data); long Test(double* data);
long Test(double* data, long result);
}; };
#endif #endif
File diff suppressed because it is too large Load Diff
+51 -44
View File
@@ -1,63 +1,70 @@
#include "DecisionTree.h" #include "DecisionTree.h"
#include "aquery.h" #include "RF.h"
// __AQ_NO_SESSION__ // __AQ_NO_SESSION__
#include "../server/table.h" #include "../server/table.h"
#include "aquery.h"
DecisionTree *dt = nullptr; DecisionTree *dt = nullptr;
RandomForest *rf = nullptr;
__AQEXPORT__(bool) newtree(int height, long f, ColRef<int> sparse, double forget, long maxf, long noclasses, Evaluation e, long r, long rb){ __AQEXPORT__(bool)
if(sparse.size!=f)return 0; newtree(int height, long f, ColRef<int> X, double forget, long maxf, long noclasses, Evaluation e, long r, long rb)
int* issparse = (int*)malloc(f*sizeof(int)); {
for(long i=0; i<f; i++){ if (X.size != f)
issparse[i] = sparse.container[i]; return false;
} int *X_cpy = (int *)malloc(f * sizeof(int));
if(maxf<0)maxf=f;
dt = new DecisionTree(height, f, issparse, forget, maxf, noclasses, e, r, rb);
return 1;
}
__AQEXPORT__(bool) fit(ColRef<ColRef<double>> X, ColRef<int> y){ memcpy(X_cpy, X.container, f);
if(X.size != y.size)return 0;
double** data = (double**)malloc(X.size*sizeof(double*)); if (maxf < 0)
long* result = (long*)malloc(y.size*sizeof(long)); maxf = f;
for(long i=0; i<X.size; i++){ dt = new DecisionTree(f, X_cpy, forget, maxf, noclasses, e);
data[i] = X.container[i].container; rf = new RandomForest(height, f, X_cpy, forget, noclasses, e)
result[i] = y.container[i];
}
data[pt] = (double*)malloc(X.size*sizeof(double));
for(j=0; j<X.size; j++){
data[pt][j]=X.container[j];
}
result[pt] = y;
pt ++;
return 1;
}
__AQEXPORT__(bool) fit(vector_type<vector_type<double>> v, vector_type<long> res){
double** data = (double**)malloc(v.size*sizeof(double*));
for(int i = 0; i < v.size; ++i)
data[i] = v.container[i].container;
dt->fit(data, res.container, v.size);
return true; return true;
} }
__AQEXPORT__(vectortype_cstorage) predict(vector_type<vector_type<double>> v){ // size_t pt = 0;
// __AQEXPORT__(bool) fit(ColRef<ColRef<double>> X, ColRef<int> y){
// if(X.size != y.size)return 0;
// double** data = (double**)malloc(X.size*sizeof(double*));
// long* result = (long*)malloc(y.size*sizeof(long));
// for(long i=0; i<X.size; i++){
// data[i] = X.container[i].container;
// result[i] = y.container[i];
// }
// data[pt] = (double*)malloc(X.size*sizeof(double));
// for(uint32_t j=0; j<X.size; j++){
// data[pt][j]=X.container[j];
// }
// result[pt] = y;
// pt ++;
// return 1;
// }
__AQEXPORT__(bool)
fit(vector_type<vector_type<double>> v, vector_type<long> res)
{
double **data = (double **)malloc(v.size * sizeof(double *));
for (int i = 0; i < v.size; ++i)
data[i] = v.container[i].container;
// dt->fit(data, res.container, v.size);
rf->fit(data, res.container, v.size);
return true;
}
__AQEXPORT__(vectortype_cstorage)
predict(vector_type<vector_type<double>> v)
{
int *result = (int *)malloc(v.size * sizeof(int)); int *result = (int *)malloc(v.size * sizeof(int));
for(long i=0; i<v.size; i++){ for (long i = 0; i < v.size; i++)
result[i]=dt->Test(v.container[i].container, dt->DTree); //result[i] = dt->Test(v.container[i].container, dt->DTree);
//printf("%d ", result[i]); result[i] = rf->Test(v.container, rf->DTrees);
}
auto container = (vector_type<int> *)malloc(sizeof(vector_type<int>)); auto container = (vector_type<int> *)malloc(sizeof(vector_type<int>));
container->size = v.size; container->size = v.size;
container->capacity = 0; container->capacity = 0;
container->container = result; container->container = result;
// container->out(10);
// ColRef<vector_type<int>>* col = (ColRef<vector_type<int>>*)malloc(sizeof(ColRef<vector_type<int>>));
auto ret = vectortype_cstorage{.container = container, .size = 1, .capacity = 0}; auto ret = vectortype_cstorage{.container = container, .size = 1, .capacity = 0};
// col->initfrom(ret, "sibal");
// print(*col);
return ret; return ret;
//return true;
} }
+8 -5
View File
@@ -31,12 +31,14 @@ void print<__int128_t>(const __int128_t& v, const char* delimiter){
s[40] = 0; s[40] = 0;
std::cout<< get_int128str(v, s+40)<< delimiter; std::cout<< get_int128str(v, s+40)<< delimiter;
} }
template <> template <>
void print<__uint128_t>(const __uint128_t&v, const char* delimiter){ void print<__uint128_t>(const __uint128_t&v, const char* delimiter){
char s[41]; char s[41];
s[40] = 0; s[40] = 0;
std::cout<< get_uint128str(v, s+40) << delimiter; std::cout<< get_uint128str(v, s+40) << delimiter;
} }
std::ostream& operator<<(std::ostream& os, __int128 & v) std::ostream& operator<<(std::ostream& os, __int128 & v)
{ {
print(v); print(v);
@@ -76,6 +78,7 @@ char* intToString(T val, char* buf){
return buf; return buf;
} }
void skip(const char*& buf){ void skip(const char*& buf){
while(*buf && (*buf >'9' || *buf < '0')) buf++; while(*buf && (*buf >'9' || *buf < '0')) buf++;
} }
@@ -264,7 +267,7 @@ std::string base62uuid(int l) {
static uniform_int_distribution<uint64_t> u(0x10000, 0xfffff); static uniform_int_distribution<uint64_t> u(0x10000, 0xfffff);
uint64_t uuid = (u(engine) << 32ull) + uint64_t uuid = (u(engine) << 32ull) +
(std::chrono::system_clock::now().time_since_epoch().count() & 0xffffffff); (std::chrono::system_clock::now().time_since_epoch().count() & 0xffffffff);
//printf("%llu\n", uuid);
string ret; string ret;
while (uuid && l-- >= 0) { while (uuid && l-- >= 0) {
ret = string("") + base62alp[uuid % 62] + ret; ret = string("") + base62alp[uuid % 62] + ret;
@@ -283,10 +286,10 @@ class A{
std::chrono::high_resolution_clock::time_point tp; std::chrono::high_resolution_clock::time_point tp;
A(){ A(){
tp = std::chrono::high_resolution_clock::now(); tp = std::chrono::high_resolution_clock::now();
printf("A %llx created.\n", tp.time_since_epoch().count()); printf("A %llu created.\n", tp.time_since_epoch().count());
} }
~A() { ~A() {
printf("A %llx died after %lldns.\n", tp.time_since_epoch().count(), printf("A %llu died after %lldns.\n", tp.time_since_epoch().count(),
(std::chrono::high_resolution_clock::now() - tp).count()); (std::chrono::high_resolution_clock::now() - tp).count());
} }
}; };
@@ -525,8 +528,8 @@ void ScratchSpace::cleanup(){
static_cast<vector_type<void*>*>(temp_memory_fractions); static_cast<vector_type<void*>*>(temp_memory_fractions);
if (vec_tmpmem_fractions->size) { if (vec_tmpmem_fractions->size) {
for(auto& mem : *vec_tmpmem_fractions){ for(auto& mem : *vec_tmpmem_fractions){
free(mem); //free(mem);
//GC::gc_handle->reg(mem); GC::gc_handle->reg(mem);
} }
vec_tmpmem_fractions->clear(); vec_tmpmem_fractions->clear();
} }
+84 -2
View File
@@ -73,16 +73,17 @@ struct StoredProcedure{
void **__rt_loaded_modules; void **__rt_loaded_modules;
}; };
struct Context { struct Context {
typedef int (*printf_type) (const char *format, ...); typedef int (*printf_type) (const char *format, ...);
void* module_function_maps = 0; void* module_function_maps = nullptr;
Config* cfg; Config* cfg;
int n_buffers, *sz_bufs; int n_buffers, *sz_bufs;
void **buffers; void **buffers;
void* alt_server = 0; void* alt_server = nullptr;
Log_level log_level = LOG_INFO; Log_level log_level = LOG_INFO;
Session current; Session current;
@@ -115,6 +116,12 @@ struct Context{
}; };
struct StoredProcedurePayload {
StoredProcedure *p;
Context* cxt;
};
int execTriggerPayload(void*);
#ifdef _WIN32 #ifdef _WIN32
#define __DLLEXPORT__ __declspec(dllexport) __stdcall #define __DLLEXPORT__ __declspec(dllexport) __stdcall
@@ -169,4 +176,79 @@ inline void AQ_ZeroMemory(_This_Struct& __val) {
memset(&__val, 0, sizeof(_This_Struct)); memset(&__val, 0, sizeof(_This_Struct));
} }
#ifdef __USE_STD_SEMAPHORE__
#include <semaphore>
class A_Semaphore {
private:
std::binary_semaphore native_handle;
public:
A_Semaphore(bool v = false) {
native_handle = std::binary_semaphore(v);
}
void acquire() {
native_handle.acquire();
}
void release() {
native_handle.release();
}
~A_Semaphore() { }
};
#else
#ifdef _WIN32
class A_Semaphore {
private:
void* native_handle;
public:
A_Semaphore(bool);
void acquire();
void release();
~A_Semaphore();
};
#else
#ifdef __APPLE__
#include <dispatch/dispatch.h>
class A_Semaphore {
private:
dispatch_semaphore_t native_handle;
public:
explicit A_Semaphore(bool v = false) {
native_handle = dispatch_semaphore_create(v);
}
void acquire() {
// puts("acquire");
dispatch_semaphore_wait(native_handle, DISPATCH_TIME_FOREVER);
}
void release() {
// puts("release");
dispatch_semaphore_signal(native_handle);
}
~A_Semaphore() {
}
};
#else
#include <semaphore.h>
class A_Semaphore {
private:
sem_t native_handle;
public:
A_Semaphore(bool v = false) {
sem_init(&native_handle, v, 1);
}
void acquire() {
sem_wait(&native_handle);
}
void release() {
sem_post(&native_handle);
}
~A_Semaphore() {
sem_destroy(&native_handle);
}
};
#endif // __APPLE__
#endif // _WIN32
#endif //__USE_STD_SEMAPHORE__
void print_monetdb_results(void* _srv, const char* sep, const char* end, uint32_t limit);
#endif #endif
+5 -5
View File
@@ -6,17 +6,17 @@
struct Context; struct Context;
struct Server{ struct Server{
MYSQL *server = 0; MYSQL *server = nullptr;
Context *cxt = 0; Context *cxt = nullptr;
bool status = 0; bool status = false;
bool has_error = false; bool has_error = false;
char* query = 0; char* query = nullptr;
int type = 0; int type = 0;
void connect(Context* cxt, const char* host = "bill.local", void connect(Context* cxt, const char* host = "bill.local",
const char* user = "root", const char* passwd = "0508", const char* user = "root", const char* passwd = "0508",
const char* db_name = "db", const unsigned int port = 3306, const char* db_name = "db", const unsigned int port = 3306,
const char* unix_socket = 0, const unsigned long client_flag = 0 const char* unix_socket = nullptr, const unsigned long client_flag = 0
); );
void exec(const char* q); void exec(const char* q);
void close(); void close();
+61 -2
View File
@@ -5,11 +5,25 @@
#include <string> #include <string>
#include "monetdb_conn.h" #include "monetdb_conn.h"
#include "monetdbe.h" #include "monetdbe.h"
#include "table.h" #include "table.h"
#include <thread> #include <thread>
#ifdef _WIN32
#include "winhelper.h"
#else
#include <dlfcn.h>
#include <fcntl.h>
#include <sys/mman.h>
#include <atomic>
#endif // _WIN32
#ifdef ERROR
#undef ERROR #undef ERROR
#endif
#ifdef static_assert
#undef static_assert #undef static_assert
#endif
constexpr const char* monetdbe_type_str[] = { constexpr const char* monetdbe_type_str[] = {
"monetdbe_bool", "monetdbe_int8_t", "monetdbe_int16_t", "monetdbe_int32_t", "monetdbe_int64_t", "monetdbe_bool", "monetdbe_int8_t", "monetdbe_int16_t", "monetdbe_int32_t", "monetdbe_int64_t",
@@ -110,7 +124,7 @@ void Server::exec(const char* q){
auto _res = static_cast<monetdbe_result*>(this->res); auto _res = static_cast<monetdbe_result*>(this->res);
monetdbe_cnt _cnt = 0; monetdbe_cnt _cnt = 0;
auto qresult = monetdbe_query(*server, const_cast<char*>(q), &_res, &_cnt); auto qresult = monetdbe_query(*server, const_cast<char*>(q), &_res, &_cnt);
if (_res != 0){ if (_res != nullptr){
this->cnt = _res->nrows; this->cnt = _res->nrows;
this->res = _res; this->res = _res;
} }
@@ -152,6 +166,8 @@ void Server::print_results(const char* sep, const char* end){
col_data[i] = static_cast<char *>(cols[i]->data); col_data[i] = static_cast<char *>(cols[i]->data);
szs [i] = monetdbe_type_szs[cols[i]->type]; szs [i] = monetdbe_type_szs[cols[i]->type];
header_string = header_string + cols[i]->name + sep + '|' + sep; header_string = header_string + cols[i]->name + sep + '|' + sep;
if (err_msg) [[unlikely]]
puts(err_msg);
} }
if (const size_t l_sep = strlen(sep) + 1; header_string.size() >= l_sep) if (const size_t l_sep = strlen(sep) + 1; header_string.size() >= l_sep)
header_string.resize(header_string.size() - l_sep); header_string.resize(header_string.size() - l_sep);
@@ -173,7 +189,7 @@ void Server::print_results(const char* sep, const char* end){
void Server::close(){ void Server::close(){
if(this->server){ if(this->server){
auto server = static_cast<monetdbe_database*>(this->server); auto server = static_cast<monetdbe_database*>(this->server);
monetdbe_close(*(server)); monetdbe_close(*server);
free(server); free(server);
this->server = nullptr; this->server = nullptr;
} }
@@ -215,3 +231,46 @@ bool Server::havehge() {
return false; return false;
#endif #endif
} }
void ExecuteStoredProcedureEx(const StoredProcedure *p, Context* cxt){
auto server = static_cast<Server*>(cxt->alt_server);
void* handle = nullptr;
uint32_t procedure_module_cursor = 0;
for(uint32_t i = 0; i < p->cnt; ++i) {
switch(p->queries[i][0]){
case 'Q': {
server->exec(p->queries[i]);
}
break;
case 'P': {
auto c = code_snippet(dlsym(handle, p->queries[i]+1));
c(cxt);
}
break;
case 'N': {
if(procedure_module_cursor < p->postproc_modules)
handle = p->__rt_loaded_modules[procedure_module_cursor++];
}
break;
case 'O': {
uint32_t limit;
memcpy(&limit, p->queries[i] + 1, sizeof(uint32_t));
if (limit == 0)
continue;
print_monetdb_results(server, " ", "\n", limit);
}
break;
default:
printf("Warning Q%u: unrecognized command %c.\n",
i, p->queries[i][0]);
}
}
}
int execTriggerPayload(void* args) {
auto spp = (StoredProcedurePayload*)(args);
ExecuteStoredProcedureEx(spp->p, spp->cxt);
delete spp;
return 0;
}
+14 -9
View File
@@ -4,18 +4,18 @@
struct Context; struct Context;
struct Server{ struct Server{
void *server = 0; void *server = nullptr;
Context *cxt = 0; Context *cxt = nullptr;
bool status = 0; bool status = false;
char* query = 0; char* query = nullptr;
int type = 1; int type = 1;
void* res = 0; void* res = nullptr;
void* ret_col = 0; void* ret_col = nullptr;
long long cnt = 0; long long cnt = 0;
char* last_error = 0; char* last_error = nullptr;
Server(Context* cxt = nullptr); explicit Server(Context* cxt = nullptr);
void connect(Context* cxt); void connect(Context* cxt);
void exec(const char* q); void exec(const char* q);
void *getCol(int col_idx); void *getCol(int col_idx);
@@ -24,7 +24,7 @@ struct Server{
static bool havehge(); static bool havehge();
void test(const char*); void test(const char*);
void print_results(const char* sep = " ", const char* end = "\n"); void print_results(const char* sep = " ", const char* end = "\n");
friend void print_monetdb_results(Server* srv, const char* sep, const char* end, int limit); friend void print_monetdb_results(void* _srv, const char* sep, const char* end, int limit);
~Server(); ~Server();
}; };
@@ -34,4 +34,9 @@ struct monetdbe_table_data{
void* cols; void* cols;
}; };
size_t
monetdbe_get_size(void* dbhdl, const char *table_name);
void*
monetdbe_get_col(void* dbhdl, const char *table_name, uint32_t col_id);
#endif #endif
+90
View File
@@ -0,0 +1,90 @@
// Non-standard Extensions for MonetDBe, may break concurrency control!
#include "monetdbe.h"
#include <stdint.h>
#include "mal_client.h"
#include "sql_mvc.h"
#include "sql_semantic.h"
#include "mal_exception.h"
typedef struct column_storage {
int refcnt;
int bid;
int ebid; /* extra bid */
int uibid; /* bat with positions of updates */
int uvbid; /* bat with values of updates */
storage_type st; /* ST_DEFAULT, ST_DICT, ST_FOR */
bool cleared;
bool merged; /* only merge changes once */
size_t ucnt; /* number of updates */
ulng ts; /* version timestamp */
} column_storage;
typedef struct segment {
BUN start;
BUN end;
bool deleted; /* we need to keep a dense segment set, 0 - end of last segemnt,
some segments maybe deleted */
ulng ts; /* timestamp on this segment, ie tid of some active transaction or commit time of append/delete or
rollback time, ie ready for reuse */
ulng oldts; /* keep previous ts, for rollbacks */
struct segment *next; /* usualy one should be enough */
struct segment *prev; /* used in destruction list */
} segment;
/* container structure to allow sharing this structure */
typedef struct segments {
sql_ref r;
struct segment *h;
struct segment *t;
} segments;
typedef struct storage {
column_storage cs; /* storage on disk */
segments *segs; /* local used segements */
struct storage *next;
} storage;
typedef struct {
char language; /* 'S' or 's' or 'X' */
char depth; /* depth >= 1 means no output for trans/schema statements */
int remote; /* counter to make remote function names unique */
mvc *mvc;
char others[];
} backend;
typedef struct {
Client c;
char *msg;
monetdbe_data_blob blob_null;
monetdbe_data_date date_null;
monetdbe_data_time time_null;
monetdbe_data_timestamp timestamp_null;
str mid;
} monetdbe_database_internal;
size_t
monetdbe_get_size(monetdbe_database dbhdl, const char *table_name)
{
monetdbe_database_internal* hdl = (monetdbe_database_internal*)dbhdl;
backend* be = ((backend *)(((monetdbe_database_internal*)dbhdl)->c->sqlcontext));
mvc *m = be->mvc;
sql_table *t = find_table_or_view_on_scope(m, NULL, "sys", table_name, "CATALOG", false);
sql_column *col = ol_first_node(t->columns)->data;
sqlstore* store = m->store;
size_t sz = store->storage_api.count_col(m->session->tr, col, QUICK);
return sz;
}
void*
monetdbe_get_col(monetdbe_database dbhdl, const char *table_name, uint32_t col_id) {
monetdbe_database_internal* hdl = (monetdbe_database_internal*)dbhdl;
backend* be = ((backend *)(((monetdbe_database_internal*)dbhdl)->c->sqlcontext));
mvc *m = be->mvc;
sql_table *t = find_table_or_view_on_scope(m, NULL, "sys", table_name, "CATALOG", false);
sql_column *col = ol_fetch(t->columns, col_id);
sqlstore* store = m->store;
BAT *b = store->storage_api.bind_col(m->session->tr, col, QUICK);
BATiter iter = bat_iterator(b);
return iter.base;
}
+1 -1
View File
@@ -6,7 +6,7 @@ template <class Comparator, typename T = uint32_t>
class priority_vector : public vector_type<T> { class priority_vector : public vector_type<T> {
const Comparator comp; const Comparator comp;
public: public:
priority_vector(Comparator comp = std::less<T>{}) : explicit priority_vector(Comparator comp = std::less<T>{}) :
comp(comp), vector_type<T>(0) {} comp(comp), vector_type<T>(0) {}
void emplace_back(T val) { void emplace_back(T val) {
vector_type<T>::emplace_back(val); vector_type<T>::emplace_back(val);
+4 -69
View File
@@ -20,10 +20,6 @@
#include <sys/mman.h> #include <sys/mman.h>
#include <atomic> #include <atomic>
// fast numeric to string conversion
#include "jeaiii_to_text.h"
#include "dragonbox/dragonbox_to_chars.h"
struct SharedMemory struct SharedMemory
{ {
std::atomic<bool> a; std::atomic<bool> a;
@@ -41,69 +37,7 @@ struct SharedMemory
} }
}; };
#ifndef __USE_STD_SEMAPHORE__ #endif // _WIN32
#ifdef __APPLE__
#include <dispatch/dispatch.h>
class A_Semaphore {
private:
dispatch_semaphore_t native_handle;
public:
A_Semaphore(bool v = false) {
native_handle = dispatch_semaphore_create(v);
}
void acquire() {
// puts("acquire");
dispatch_semaphore_wait(native_handle, DISPATCH_TIME_FOREVER);
}
void release() {
// puts("release");
dispatch_semaphore_signal(native_handle);
}
~A_Semaphore() {
}
};
#else
#include <semaphore.h>
class A_Semaphore {
private:
sem_t native_handle;
public:
A_Semaphore(bool v = false) {
sem_init(&native_handle, v, 1);
}
void acquire() {
sem_wait(&native_handle);
}
void release() {
sem_post(&native_handle);
}
~A_Semaphore() {
sem_destroy(&native_handle);
}
};
#endif
#endif
#endif
#ifdef __USE_STD_SEMAPHORE__
#define __AQUERY_ITC_USE_SEMPH__
#include <semaphore>
class A_Semaphore {
private:
std::binary_semaphore native_handle;
public:
A_Semaphore(bool v = false) {
native_handle = std::binary_semaphore(v);
}
void acquire() {
native_handle.acquire();
}
void release() {
native_handle.release();
}
~A_Semaphore() { }
};
#endif
#ifdef __AQUERY_ITC_USE_SEMPH__ #ifdef __AQUERY_ITC_USE_SEMPH__
A_Semaphore prompt{ true }, engine{ false }; A_Semaphore prompt{ true }, engine{ false };
@@ -225,8 +159,9 @@ inline constexpr static unsigned char monetdbe_type_szs[] = {
1 1
}; };
constexpr uint32_t output_buffer_size = 65536; constexpr uint32_t output_buffer_size = 65536;
void print_monetdb_results(Server* srv, const char* sep = " ", const char* end = "\n", void print_monetdb_results(void* _srv, const char* sep = " ", const char* end = "\n",
uint32_t limit = std::numeric_limits<uint32_t>::max()) { uint32_t limit = std::numeric_limits<uint32_t>::max()) {
auto srv = static_cast<Server *>(_srv);
if (!srv->haserror() && srv->cnt && limit) { if (!srv->haserror() && srv->cnt && limit) {
char buffer[output_buffer_size]; char buffer[output_buffer_size];
auto _res = static_cast<monetdbe_result*> (srv->res); auto _res = static_cast<monetdbe_result*> (srv->res);
@@ -255,7 +190,7 @@ void print_monetdb_results(Server* srv, const char* sep = " ", const char* end =
puts("Error: separator or end string too long"); puts("Error: separator or end string too long");
goto cleanup; goto cleanup;
} }
if (header_string.size() - l_sep - 1>= 0) if (header_string.size() >= l_sep + 1)
header_string.resize(header_string.size() - l_sep - 1); header_string.resize(header_string.size() - l_sep - 1);
header_string += end + std::string(header_string.size(), '=') + end; header_string += end + std::string(header_string.size(), '=') + end;
fputs(header_string.c_str(), stdout); fputs(header_string.c_str(), stdout);
+26 -16
View File
@@ -4,21 +4,26 @@
#include <thread> #include <thread>
#include <cstdio> #include <cstdio>
#include <cstdlib> #include <cstdlib>
using namespace std;
FILE *fp; FILE *fp;
long long testing_throughput(uint32_t n_jobs, bool prompt = true){ long long testing_throughput(uint32_t n_jobs, bool prompt = true){
using namespace std::chrono_literals;
printf("Threadpool througput test with %u jobs. Press any key to start.\n", n_jobs); printf("Threadpool througput test with %u jobs. Press any key to start.\n", n_jobs);
auto tp = ThreadPool(thread::hardware_concurrency()); auto tp = ThreadPool(std::thread::hardware_concurrency());
getchar(); getchar();
auto i = 0u; auto i = 0u;
fp = fopen("tmp.tmp", "wb"); fp = fopen("tmp.tmp", "wb");
auto time = chrono::high_resolution_clock::now(); auto time = std::chrono::high_resolution_clock::now();
while(i++ < n_jobs) tp.enqueue_task({ [](void* f) {fprintf(fp, "%d ", *(int*)f); free(f); }, new int(i) }); while(i++ < n_jobs) {
payload_t payload;
payload.f = [](void* f) {fprintf(fp, "%d ", *(int*)f); free(f); return 0; };
payload.args = new int(i);
tp.enqueue_task(payload);
}
puts("done dispatching."); puts("done dispatching.");
while (tp.busy()) this_thread::sleep_for(1s); while (tp.busy()) std::this_thread::sleep_for(1s);
auto t = (chrono::high_resolution_clock::now() - time).count(); auto t = (std::chrono::high_resolution_clock::now() - time).count();
printf("\nTr: %u, Ti: %lld \nThroughput: %lf transactions/ns\n", i, t, i/(double)(t)); printf("\nTr: %u, Ti: %lld \nThroughput: %lf transactions/ns\n", i, t, i/(double)(t));
//this_thread::sleep_for(2s); //this_thread::sleep_for(2s);
fclose(fp); fclose(fp);
@@ -27,26 +32,31 @@ long long testing_throughput(uint32_t n_jobs, bool prompt = true){
long long testing_transaction(uint32_t n_burst, uint32_t n_batch, long long testing_transaction(uint32_t n_burst, uint32_t n_batch,
uint32_t base_time, uint32_t var_time, bool prompt = true, FILE* _fp = stdout){ uint32_t base_time, uint32_t var_time, bool prompt = true, FILE* _fp = stdout){
using namespace std::chrono_literals;
printf("Threadpool transaction test: burst: %u, batch: %u, time: [%u, %u].\n" printf("Threadpool transaction test: burst: %u, batch: %u, time: [%u, %u].\n"
, n_burst, n_batch, base_time, var_time + base_time); , n_burst, n_batch, base_time, var_time + base_time);
if (prompt) { if (prompt) {
puts("Press any key to start."); puts("Press any key to start.");
getchar(); getchar();
} }
auto tp = ThreadPool(thread::hardware_concurrency()); auto tp = ThreadPool(std::thread::hardware_concurrency());
fp = _fp; fp = _fp;
auto i = 0u, j = 0u; auto i = 0u, j = 0u;
auto time = chrono::high_resolution_clock::now(); auto time = std::chrono::high_resolution_clock::now();
while(j++ < n_batch){ while(j++ < n_batch){
i = 0u; i = 0u;
while(i++ < n_burst) while(i++ < n_burst) {
tp.enqueue_task({ [](void* f) { fprintf(fp, "%d ", *(int*)f); free(f); }, new int(j) }); payload_t payload;
payload.f = [](void* f) { fprintf(fp, "%d ", *(int*)f); free(f); return 0; };
payload.args = new int(j);
tp.enqueue_task(payload);
}
fflush(stdout); fflush(stdout);
this_thread::sleep_for(chrono::microseconds(rand()%var_time + base_time)); std::this_thread::sleep_for(std::chrono::microseconds(rand()%var_time + base_time));
} }
puts("done dispatching."); puts("done dispatching.");
while (tp.busy()) this_thread::sleep_for(1s); while (tp.busy()) std::this_thread::sleep_for(1s);
auto t = (chrono::high_resolution_clock::now() - time).count(); auto t = (std::chrono::high_resolution_clock::now() - time).count();
printf("\nTr: %u, Ti: %lld \nThroughput: %lf transactions/ns\n", j*i, t, j*i/(double)(t)); printf("\nTr: %u, Ti: %lld \nThroughput: %lf transactions/ns\n", j*i, t, j*i/(double)(t));
return t; return t;
@@ -58,17 +68,17 @@ long long testing_destruction(bool prompt = true){
puts("Press any key to start."); puts("Press any key to start.");
getchar(); getchar();
} }
auto time = chrono::high_resolution_clock::now(); auto time = std::chrono::high_resolution_clock::now();
for(int i = 0; i < 8; ++i) for(int i = 0; i < 8; ++i)
testing_transaction(0xfff, 0xff, 400, 100, false, fp); testing_transaction(0xfff, 0xff, 400, 100, false, fp);
for(int i = 0; i < 64; ++i) for(int i = 0; i < 64; ++i)
testing_transaction(0xff, 0xf, 60, 20, false, fp); testing_transaction(0xff, 0xf, 60, 20, false, fp);
for(int i = 0; i < 1024; ++i) { for(int i = 0; i < 1024; ++i) {
auto tp = new ThreadPool(256); auto tp = new ThreadPool(255);
delete tp; delete tp;
} }
return 0; return 0;
auto t = (chrono::high_resolution_clock::now() - time).count(); auto t = (std::chrono::high_resolution_clock::now() - time).count();
fclose(fp); fclose(fp);
return t; return t;
} }
+44 -13
View File
@@ -1,4 +1,5 @@
#include "threading.h" #include "threading.h"
#include "libaquery.h"
#include <thread> #include <thread>
#include <atomic> #include <atomic>
#include <mutex> #include <mutex>
@@ -105,7 +106,8 @@ ThreadPool::ThreadPool(uint32_t n_threads)
current_payload = new payload_t[n_threads]; current_payload = new payload_t[n_threads];
for (uint32_t i = 0; i < n_threads; ++i){ for (uint32_t i = 0; i < n_threads; ++i){
atomic_init(tf + i, static_cast<unsigned char>(0b10)); // atomic_init(tf + i, static_cast<unsigned char>(0b10));
tf[i] = static_cast<unsigned char>(0b10);
th[i] = thread(&ThreadPool::daemon_proc, this, i); th[i] = thread(&ThreadPool::daemon_proc, this, i);
} }
@@ -152,21 +154,50 @@ bool ThreadPool::busy(){
return true; return true;
} }
Trigger::Trigger(ThreadPool* tp){ IntervalBasedTriggerHost::IntervalBasedTriggerHost(ThreadPool* tp){
this->tp = tp;
this->triggers = new vector_type<IntervalBasedTrigger>;
trigger_queue_lock = new mutex();
this->now = std::chrono::high_resolution_clock::now().time_since_epoch().count();
} }
void IntervalBasedTrigger::timer::reset(){ void IntervalBasedTriggerHost::add_trigger(StoredProcedure *p, uint32_t interval) {
auto tr = IntervalBasedTrigger{.interval = interval, .time_remaining = 0, .sp = p};
auto vt_triggers = static_cast<vector_type<IntervalBasedTrigger> *>(this->triggers);
trigger_queue_lock->lock();
vt_triggers->emplace_back(tr);
trigger_queue_lock->unlock();
}
void IntervalBasedTriggerHost::tick() {
const auto current_time = std::chrono::high_resolution_clock::now().time_since_epoch().count();
const auto delta_t = static_cast<uint32_t>((current_time - now) / 1000000); // miliseconds precision
now = current_time;
auto vt_triggers = static_cast<vector_type<IntervalBasedTrigger> *>(this->triggers);
trigger_queue_lock->lock();
for(auto& t : *vt_triggers) {
if(t.tick(delta_t)) {
payload_t payload;
payload.f = execTriggerPayload;
payload.args = static_cast<void*>(new StoredProcedurePayload {t.sp, cxt});
tp->enqueue_task(payload);
}
}
trigger_queue_lock->unlock();
}
void IntervalBasedTrigger::reset() {
time_remaining = interval; time_remaining = interval;
} }
bool IntervalBasedTrigger::timer::tick(uint32_t t){ bool IntervalBasedTrigger::tick(uint32_t delta_t) {
if (time_remaining > t) { bool ret = false;
time_remaining -= t; if (time_remaining <= delta_t)
return false; ret = true;
} if (auto curr_dt = delta_t % interval; time_remaining <= curr_dt)
else{ time_remaining = interval + time_remaining - curr_dt;
time_remaining = interval - t%interval; else
return true; time_remaining = time_remaining - curr_dt;
} return ret;
} }
+27 -11
View File
@@ -19,7 +19,7 @@ struct payload_t{
class ThreadPool{ class ThreadPool{
public: public:
ThreadPool(uint32_t n_threads = 0); explicit ThreadPool(uint32_t n_threads = 0);
void enqueue_task(const payload_t& payload); void enqueue_task(const payload_t& payload);
bool busy(); bool busy();
virtual ~ThreadPool(); virtual ~ThreadPool();
@@ -39,29 +39,45 @@ private:
}; };
class Trigger{ #include <thread>
private: #include <mutex>
void* triggers; //min-heap by t-rem class A_Semphore;
virtual void tick() = 0;
class TriggerHost {
protected:
void* triggers;
std::thread* handle;
ThreadPool *tp;
Context* cxt;
std::mutex* trigger_queue_lock;
virtual void tick() = 0;
public: public:
Trigger(ThreadPool* tp); TriggerHost() = default;
virtual ~TriggerHost() = default;
}; };
class IntervalBasedTrigger : public Trigger{ struct StoredProcedure;
public:
struct timer{ struct IntervalBasedTrigger {
uint32_t interval; // in milliseconds uint32_t interval; // in milliseconds
uint32_t time_remaining; uint32_t time_remaining;
StoredProcedure* sp;
void reset(); void reset();
bool tick(uint32_t t); bool tick(uint32_t t);
}; };
void add_trigger();
class IntervalBasedTriggerHost : public TriggerHost {
public:
explicit IntervalBasedTriggerHost(ThreadPool *tp);
void add_trigger(StoredProcedure* stored_procedure, uint32_t interval);
void remove_trigger(uint32_t tid);
private: private:
unsigned long long now;
void tick() override; void tick() override;
}; };
class CallbackBasedTrigger : public Trigger{ class CallbackBasedTriggerHost : public TriggerHost {
public: public:
void add_trigger(); void add_trigger();
private: private:
-12
View File
@@ -16,18 +16,6 @@ struct SharedMemory
void FreeMemoryMap(); void FreeMemoryMap();
}; };
#ifndef __USE_STD_SEMAPHORE__
class A_Semaphore {
private:
void* native_handle;
public:
A_Semaphore(bool);
void acquire();
void release();
~A_Semaphore();
};
#endif
#endif // WIN32 #endif // WIN32
#endif // WINHELPER #endif // WINHELPER
+2
View File
@@ -6,6 +6,8 @@ LOAD MODULE FROM "./libirf.so"
); );
create table source(x1 double, x2 double, x3 double, x4 double, x5 int64); create table source(x1 double, x2 double, x3 double, x4 double, x5 int64);
-- Create trigger 1 ~~ to predict whenever sz(source > ?)
-- Create trigger 2 ~~ to auto feed ~
load data infile "data/benchmark" into table source fields terminated by ","; load data infile "data/benchmark" into table source fields terminated by ",";
create table sparse(x int); create table sparse(x int);