Changeset: 1cc40e8acd9f for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=1cc40e8acd9f
Modified Files:
clients/mapiclient/mclient.c
clients/mapilib/Makefile.ag
clients/mapilib/mapi.c
common/stream/stream.h
common/utils/mcrypt.c
configure.ag
monetdb5/modules/mal/mal_mapi.c
sql/backends/monet5/Makefile.ag
sql/backends/monet5/sql_result.c
Branch: protocol
Log Message:
Add TurboPFOR and Binpack column compression techniques, and optimize them so
they write directly to stream buffer (to avoid unnecessary memcpy).
diffs (truncated from 463 to 300 lines):
diff --git a/clients/mapiclient/mclient.c b/clients/mapiclient/mclient.c
--- a/clients/mapiclient/mclient.c
+++ b/clients/mapiclient/mclient.c
@@ -2957,7 +2957,7 @@ usage(const char *prog, int xit)
fprintf(stderr, " -C version | --compression=type specify compression
method {snappy,lz4}\n");
fprintf(stderr, " -P version | --protocol=version specify protocol
version {prot9,prot10,prot10compressed}\n");
fprintf(stderr, " -B size | --blocksize=size specify protocol
block size (>= %d)\n", BLOCK);
- fprintf(stderr, " -c colcomp | --colcomp=type specify column
compression type {none,pfor,protobuf}");
+ fprintf(stderr, " -c colcomp | --colcomp=type specify column
compression type {none,pfor,protobuf,binpack}");
fprintf(stderr, " -H | --history load/save cmdline
history (default off)\n");
fprintf(stderr, " -i | --interactive[=tm] interpret `\\'
commands on stdin, use time formatting {ms,s,m}\n");
diff --git a/clients/mapilib/Makefile.ag b/clients/mapilib/Makefile.ag
--- a/clients/mapilib/Makefile.ag
+++ b/clients/mapilib/Makefile.ag
@@ -7,7 +7,7 @@
MTSAFE
INCLUDES = ../../common/options ../../common/stream ../../common/utils \
- $(MSGCONTROL_FLAGS) $(pfor_CFLAGS)
+ $(MSGCONTROL_FLAGS) $(pfor_CFLAGS) $(binpack_CFLAGS)
lib_mapi = {
VERSION = $(MAPI_VERSION)
@@ -15,7 +15,7 @@ lib_mapi = {
LIBS = $(SOCKET_LIBS) ../../common/stream/libstream \
../../common/stream/libprotobuf \
../../common/options/libmoptions \
- ../../common/utils/libmcrypt $(openssl_LIBS) $(pfor_LIBS)
$(protobuf_LIBS)
+ ../../common/utils/libmcrypt $(openssl_LIBS) $(pfor_LIBS)
$(protobuf_LIBS) $(binpack_LIBS)
}
headers_mapi = {
diff --git a/clients/mapilib/mapi.c b/clients/mapilib/mapi.c
--- a/clients/mapilib/mapi.c
+++ b/clients/mapilib/mapi.c
@@ -812,8 +812,12 @@
# endif
#endif
+#ifdef HAVE_BINPACK
+#include <simdcomp.h>
+#endif
#ifdef HAVE_PFOR
-#include <simdcomp.h>
+#include <vint.h>
+#include <vp4dd.h>
#endif
#ifndef INVALID_SOCKET
@@ -2643,6 +2647,8 @@ mapi_reconnect(Mapi mid)
mid->colcomp = COLUMN_COMPRESSION_NONE;
} else if (strcasecmp(env_colcomp, "pfor") == 0) {
mid->colcomp = COLUMN_COMPRESSION_PFOR;
+ } else if (strcasecmp(env_colcomp, "binpack") == 0) {
+ mid->colcomp = COLUMN_COMPRESSION_BINPACK;
} else if (strcasecmp(env_colcomp, "protobuf") == 0) {
mid->colcomp = COLUMN_COMPRESSION_PROTOBUF;
}
@@ -2716,6 +2722,23 @@ mapi_reconnect(Mapi mid)
exit(1);
}
#endif
+#ifdef HAVE_BINPACK
+ if (strstr(hashes, "BINPACK")) {
+ if (mid->colcomp == COLUMN_COMPRESSION_AUTO) {
+ mid->colcomp = COLUMN_COMPRESSION_BINPACK;
+ }
+ } else if (mid->colcomp == COLUMN_COMPRESSION_BINPACK) {
+ mapi_setError(mid, "Client wants BINPACK but server
does not support it",
+ "mapi_reconnect", MERROR);
+ close_connection(mid);
+ return mid->error;
+ }
+#else
+ if (mid->colcomp == COLUMN_COMPRESSION_BINPACK) {
+ fprintf(stderr, "Client does not support BINPACK
compression.\n");
+ exit(1);
+ }
+#endif
#ifdef HAVE_LIBSNAPPY
if (strstr(hashes, "PROT10COMPR")) {
// both server and client support compressed protocol
10; use compressed version
@@ -2852,7 +2875,7 @@ mapi_reconnect(Mapi mid)
mid->database == NULL ? "" : mid->database,
prot_version == prot10 ? "PROT10" :
"PROT10COMPR",
comp == COMPRESSION_SNAPPY ? "SNAPPY" :
(comp == COMPRESSION_LZ4 ? "LZ4" : ""),
- mid->colcomp == COLUMN_COMPRESSION_PFOR ?
",HAVEPFOR" : (mid->colcomp == COLUMN_COMPRESSION_PROTOBUF ? ",PROTOBUF" :
""),
+ mid->colcomp == COLUMN_COMPRESSION_PFOR ?
",HAVEPFOR" : (mid->colcomp == COLUMN_COMPRESSION_BINPACK ? ",HAVEBINPACK" :
(mid->colcomp == COLUMN_COMPRESSION_PROTOBUF ? ",PROTOBUF" : "")),
mid->blocksize);
} else {
retval = snprintf(buf, BLOCK,
"%s:%s:%s:%s:%s:\n",
@@ -4329,7 +4352,7 @@ read_into_cache(MapiHdl hdl, int lookahe
result->fieldcnt = nr_cols;
result->maxfields = (int) nr_cols;
result->row_count = nr_rows;
- result->fields = malloc(sizeof(struct MapiColumn) *
result->fieldcnt);
+ result->fields = calloc(result->fieldcnt, sizeof(struct
MapiColumn));
result->tableid = result_set_id;
result->querytype = Q_TABLE;
result->tuple_count = 0;
@@ -5754,23 +5777,52 @@ mapi_fetch_row(MapiHdl hdl)
result->fields[i].buffer_ptr +=
sizeof(lng);
buf += col_len + sizeof(lng);
} else {
-#ifdef HAVE_PFOR
- if (hdl->mid->colcomp ==
COLUMN_COMPRESSION_PFOR && strcasecmp(result->fields[i].columntype, "int") ==
0) {
+#ifdef HAVE_BINPACK
+ if (hdl->mid->colcomp ==
COLUMN_COMPRESSION_BINPACK && strcasecmp(result->fields[i].columntype, "int")
== 0) {
lng b = *((lng*) buf);
buf += sizeof(lng);
lng length = *((lng*)(buf));
buf += sizeof(lng);
-
+ // FIXME resbuffer is not freed
uint8_t *resbuffer =
malloc(nrows * sizeof(int));
+ if (!resbuffer) {
+ return 0;
+ }
simdunpack_length((const
__m128i *)buf, nrows, (uint32_t*) resbuffer, b);
result->fields[i].buffer_ptr =
resbuffer;
buf += length;
} else {
#endif
+#ifdef HAVE_PFOR
+ if (hdl->mid->colcomp ==
COLUMN_COMPRESSION_PFOR && strcasecmp(result->fields[i].columntype, "int") ==
0) {
+ size_t n = nrows;
+ // FIXME resbuffer is not freed
+ char *resbuffer = malloc(nrows
* sizeof(int));
+ char *bufpos = resbuffer;
+ if (!resbuffer) {
+ return 0;
+ }
+ while(n > 0) {
+ size_t elements = n >
128 ? 128 : n;
+ if (elements < 128) {
+ memcpy(bufpos,
buf, elements * sizeof(int));
+ buf += elements
* sizeof(int);
+ } else {
+ buf =
p4ddecv32(buf, elements, bufpos);
+ }
+ bufpos += elements *
sizeof(int);
+ n -= elements;
+ }
+ result->fields[i].buffer_ptr =
resbuffer;
+ } else {
+#endif
buf += nrows *
result->fields[i].columnlength;
#ifdef HAVE_PFOR
}
#endif
+#ifdef HAVE_BINPACK
+ }
+#endif
}
}
result->tuple_count += nrows;
@@ -6196,12 +6248,13 @@ MapiMsg
mapi_set_column_compression(Mapi mid, const char* colcomp) {
if (strcasecmp(colcomp, "pfor") == 0) {
mid->colcomp = COLUMN_COMPRESSION_PFOR;
- }
- else if (strcasecmp(colcomp, "none") == 0) {
+ } else if (strcasecmp(colcomp, "none") == 0) {
mid->colcomp = COLUMN_COMPRESSION_NONE;
} else if (strcasecmp(colcomp, "protobuf") == 0) {
mid->colcomp = COLUMN_COMPRESSION_PROTOBUF;
- } else {
+ } else if (strcasecmp(colcomp, "binpack") == 0) {
+ mid->colcomp = COLUMN_COMPRESSION_BINPACK;
+ }else {
mapi_setError(mid, "invalid column compression type",
"mapi_set_compression", MERROR);
return -1;
}
diff --git a/common/stream/stream.h b/common/stream/stream.h
--- a/common/stream/stream.h
+++ b/common/stream/stream.h
@@ -253,7 +253,8 @@ typedef enum {
COLUMN_COMPRESSION_AUTO = 255,
COLUMN_COMPRESSION_NONE = 0,
COLUMN_COMPRESSION_PFOR = 1,
- COLUMN_COMPRESSION_PROTOBUF = 2
+ COLUMN_COMPRESSION_BINPACK = 2,
+ COLUMN_COMPRESSION_PROTOBUF = 3
} column_compression;
stream_export stream *block_stream2(stream *s, size_t bufsiz,
compression_method comp, column_compression colcomp);
diff --git a/common/utils/mcrypt.c b/common/utils/mcrypt.c
--- a/common/utils/mcrypt.c
+++ b/common/utils/mcrypt.c
@@ -43,6 +43,10 @@ mcrypt_getHashAlgorithms(void)
// the server supports PFOR
",PFOR"
#endif
+#ifdef HAVE_BINPACK
+// the server supports bin packing
+ ",BINPACK"
+#endif
);
}
diff --git a/configure.ag b/configure.ag
--- a/configure.ag
+++ b/configure.ag
@@ -1688,6 +1688,49 @@ AC_SUBST([lz4_CFLAGS])
AC_SUBST([lz4_LIBS])
AM_CONDITIONAL([HAVE_LIBLZ4], [test x$have_lz4 != xno])
+dnl check for simd binpack (de)compression library
+org_have_binpack=auto
+have_binpack=$org_have_binpack
+binpack_CFLAGS=""
+binpack_LIBS=""
+AC_ARG_WITH([binpack],
+ [AS_HELP_STRING([--with-binpack=DIR],
+ [binpack library is installed in DIR])],
+ [have_binpack="$withval"])
+
+AS_CASE(["$have_binpack"],
+ [yes|no|auto], [],
+ [
+ binpack_CFLAGS="-I$withval/include"
+ binpack_LIBS="-L$withval"])
+
+AC_MSG_CHECKING([for binpack])
+AS_VAR_IF([have_binpack], [no], [], [
+ save_CPPFLAGS="$CPPFLAGS"
+ CPPFLAGS="$CPPFLAGS $binpack_CFLAGS"
+ save_LIBS="$LIBS"
+ LIBS="$LIBS $binpack_LIBS -lsimdcomp"
+ AC_LINK_IFELSE([AC_LANG_PROGRAM([[@%:@include <stdio.h>
+@%:@include <simdcomputil.h>]], [[(void)bits(42);]])],
+ binpack_LIBS="$binpack_LIBS -lsimdcomp",
+ [ AS_VAR_IF([have_binpack], [auto], [], [
+ AC_MSG_ERROR([binpack library not found])])
+ have_binpack=no
+ why_have_binpack="(binpack library not found)" ])
+ LIBS="$save_LIBS"
+ CPPFLAGS="$save_CPPFLAGS"])
+
+AS_VAR_IF([have_binpack], [no], [], [
+ AC_DEFINE([HAVE_BINPACK], 1, [Define if you have the simd binpack
library])
+ AC_MSG_RESULT([yes: $binpack_LIBS])], [
+ binpack_CFLAGS=""
+ binpack_LIBS=""
+ AC_MSG_RESULT([no])])
+
+AC_SUBST([binpack_CFLAGS])
+AC_SUBST([binpack_LIBS])
+AM_CONDITIONAL([HAVE_BINPACK], [test x$have_binpack != xno])
+
dnl check for pfor (de)compression library
org_have_pfor=auto
@@ -1702,18 +1745,20 @@ AC_ARG_WITH([pfor],
AS_CASE(["$have_pfor"],
[yes|no|auto], [],
[
- pfor_CFLAGS="-I$withval/include"
- pfor_LIBS="-L$withval"])
+ pfor_CFLAGS="-I$withval"
+ pfor_LIBS="-Wl,$withval/libpfor.so"])
AC_MSG_CHECKING([for pfor])
AS_VAR_IF([have_pfor], [no], [], [
save_CPPFLAGS="$CPPFLAGS"
CPPFLAGS="$CPPFLAGS $pfor_CFLAGS"
save_LIBS="$LIBS"
- LIBS="$LIBS $pfor_LIBS -lsimdcomp"
+ LIBS="$LIBS $pfor_LIBS"
AC_LINK_IFELSE([AC_LANG_PROGRAM([[@%:@include <stdio.h>
-@%:@include <simdcomputil.h>]], [[(void)bits(42);]])],
- pfor_LIBS="$pfor_LIBS -lsimdcomp",
+ @%:@include <vint.h>
+ @%:@include <vp4dc.h>
+ @%:@include <stdlib.h>]], [[int *ints = malloc(sizeof(int) * 10); char
*buf = malloc(1000); char *endptr = p4denc32(ints, 10, buf);]])],
+ pfor_LIBS="$pfor_LIBS",
[ AS_VAR_IF([have_pfor], [auto], [], [
AC_MSG_ERROR([pfor library not found])])
have_pfor=no
@@ -3401,6 +3446,7 @@ echo
echo "* Available features/extensions:"
for comp in \
'atomic_ops ' \
+ 'binpack ' \
'bz2 ' \
'curl ' \
'fits ' \
@@ -3419,6 +3465,7 @@ for comp in \
'openssl ' \
'pcre ' \
'perl ' \
+ 'pfor ' \
'proj ' \
'protobuf ' \
'pthread ' \
diff --git a/monetdb5/modules/mal/mal_mapi.c b/monetdb5/modules/mal/mal_mapi.c
--- a/monetdb5/modules/mal/mal_mapi.c
+++ b/monetdb5/modules/mal/mal_mapi.c
@@ -190,9 +190,13 @@ doChallenge(void *data)
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list