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

Reply via email to