Changeset: b868df3dcc4e for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/b868df3dcc4e
Added Files:
        sql/backends/monet5/rel_copy.c
        sql/backends/monet5/rel_copy.h
Modified Files:
        sql/backends/monet5/CMakeLists.txt
        sql/backends/monet5/rel_bin.c
        sql/backends/monet5/rel_bin.h
Branch: copyparpipe
Log Message:

Move parallel COPY INTO plan generation to separate file


diffs (truncated from 688 to 300 lines):

diff --git a/sql/backends/monet5/CMakeLists.txt 
b/sql/backends/monet5/CMakeLists.txt
--- a/sql/backends/monet5/CMakeLists.txt
+++ b/sql/backends/monet5/CMakeLists.txt
@@ -130,6 +130,7 @@ target_sources(sql
   sql_assert.c sql_assert.h
   sql_upgrades.c sql_upgrades.h
   rel_bin.c rel_bin.h
+  rel_copy.c rel_copy.h
   rel_predicates.c rel_predicates.h
   sql_cat.c sql_cat.h
   sql_transaction.c sql_transaction.h
diff --git a/sql/backends/monet5/rel_bin.c b/sql/backends/monet5/rel_bin.c
--- a/sql/backends/monet5/rel_bin.c
+++ b/sql/backends/monet5/rel_bin.c
@@ -9,6 +9,7 @@
 #include "monetdb_config.h"
 
 #include "rel_bin.h"
+#include "rel_copy.h"
 #include "rel_rel.h"
 #include "rel_basetable.h"
 #include "rel_exp.h"
@@ -23,7 +24,6 @@
 #include "mal_builder.h"
 #include "opt_prelude.h"
 
-static stmt * exp_bin(backend *be, sql_exp *e, stmt *left, stmt *right, stmt 
*grp, stmt *ext, stmt *cnt, stmt *sel, int depth, int reduce, int push);
 static stmt * rel_bin(backend *be, sql_rel *rel);
 static stmt * subrel_bin(backend *be, sql_rel *rel, list *refs);
 
@@ -38,7 +38,7 @@ clean_mal_statements(backend *be, int ol
        be->mvc->errstr[0] = '\0';
 }
 
-static void
+void
 add_to_rowcount_accumulator(backend *be, int nr)
 {
        if (be->silent)
@@ -4439,299 +4439,6 @@ can_use_directappend(sql_rel *rel)
        return copy_from;
 }
 
-static int
-extract_parameter(backend *be, list *stmts, sql_exp *copyfrom, int argno)
-{
-       list *args = copyfrom->l;
-       node *n = args->h;
-       for (int i = 0; i < argno; i++)
-               n = n->next;
-       sql_exp *exp = n->data;
-       stmt *st = exp_bin(be, exp, NULL, NULL, NULL, NULL, NULL, NULL, 0, 0, 
0);
-       list_append(stmts, st);
-       return st->nr;
-}
-
-static int
-emit_receive(MalBlkPtr mb, int var_channel, int tpe)
-{
-       InstrPtr q = newAssignment(mb);
-       q = pushReturn(mb, q, var_channel);
-       q = pushArgument(mb, q, var_channel);
-       q = pushNil(mb, q, tpe);
-       return getDestVar(q);
-}
-
-static void
-emit_send(MalBlkPtr mb, int var_channel, int tpe, int var_msg)
-{
-       InstrPtr q = newAssignment(mb);
-       setReturnArgument(q, var_channel);
-       q = pushReturn(mb, q, var_msg);
-       q = pushArgument(mb, q, var_msg);
-       q = pushNil(mb, q, tpe);
-}
-
-static stmt *
-rel2bin_copyparpipe(backend *be, sql_rel *rel, list *refs, sql_exp *copyfrom)
-{
-       (void)rel;
-       (void)refs;
-       const int block_size = 1024 * 1024;
-       const int margin = 8 * 1024;
-
-       InstrPtr q;
-       MalBlkPtr mb = be->mb;
-       mvc *mvc = be->mvc;
-       sql_allocator *sa = mvc->sa;
-       list *intermediate_stmts = sa_list(sa);
-
-       int streams_type = ATOMindex("streams");
-       int bte_bat_type = newBatType(TYPE_bte);
-       int int_bat_type = newBatType(TYPE_int);
-
-       // Extract table name
-       list *copyfrom_args = copyfrom->l;
-       node *n = copyfrom_args->h;
-       sql_exp *first_arg_exp = n->data;
-       if (first_arg_exp->type != e_atom)
-               return NULL;
-       atom *first_arg_atom = first_arg_exp->l;
-       sql_table *table = first_arg_atom->data.val.pval;
-       const char *table_name = table->base.name;
-       const char *schema_name = table->s->base.name;
-       int column_count = ol_length(table->columns);
-
-       // Extract other arguments
-       int var_col_sep = extract_parameter(be, intermediate_stmts, copyfrom, 
1);
-       int var_line_sep = extract_parameter(be, intermediate_stmts, copyfrom, 
2);
-       int var_quote_char = extract_parameter(be, intermediate_stmts, 
copyfrom, 3);
-       int var_null_representation = extract_parameter(be, intermediate_stmts, 
copyfrom, 4);
-       int var_fname = extract_parameter(be, intermediate_stmts, copyfrom, 5);
-       int var_num_rows = extract_parameter(be, intermediate_stmts, copyfrom, 
6);
-       int var_offset = extract_parameter(be, intermediate_stmts, copyfrom, 7);
-       int var_best_effort = extract_parameter(be, intermediate_stmts, 
copyfrom, 8);
-       int var_fixed_width = extract_parameter(be, intermediate_stmts, 
copyfrom, 9);
-       int var_on_client = extract_parameter(be, intermediate_stmts, copyfrom, 
10);
-       int var_escape = extract_parameter(be, intermediate_stmts, copyfrom, 
11);
-
-       // coerce var_escape to bit
-       q = newStmt(mb, "calc", "!=");
-       q = pushArgument(mb, q, var_escape);
-       q = pushInt(mb, q, 0);
-       var_escape = getDestVar(q);
-
-
-       // TODO: Deal with the following
-       (void)var_num_rows;
-       (void)var_offset;
-       (void)var_best_effort;
-       (void)var_on_client;
-       (void)var_fixed_width;
-
-       q = newAssignment(mb);
-       q = pushLng(mb, q, 0);
-       int var_total_row_count = getDestVar(q);
-
-       q = newStmt(mb, "streams", "openRead");
-       q = pushArgument(mb, q, var_fname);
-       int var_stream_channel = getDestVar(q);
-
-       q = newStmt(mb, "bat", "new");
-       q = pushNil(mb, q, TYPE_bte);
-       q = pushLng(mb, q, block_size + margin);
-       int var_block_channel = getDestVar(q);
-
-       q = newAssignment(mb);
-       q = pushInt(mb, q, 0);
-       int var_skip_amounts_channel = getDestVar(q);
-
-       q = newAssignment(mb);
-       q = pushNil(mb, q, TYPE_bit);
-       int var_claim_channel = getDestVar(q);
-
-
-       // START LOOP
-       q = newAssignment(mb);
-       q->barrier = BARRIERsymbol;
-       q = pushBit(mb, q, true);
-       int var_loop_barrier = getDestVar(q);
-
-       int var_s = emit_receive(mb, var_stream_channel, streams_type);
-
-       q = newStmt(mb, "bat", "new");
-       q = pushNil(mb, q, TYPE_bte);
-       q = pushLng(mb, q, 300);
-       int var_next_block = getDestVar(q);
-
-       // START READ BLOCK
-       q = newStmt(mb, "calc", "isnotnil");
-       q->barrier = BARRIERsymbol;
-       q = pushArgument(mb, q, var_s);
-       int var_read_barrier = getDestVar(q);
-
-       q = newStmt(mb, "copy", "read");
-       q = pushArgument(mb, q, var_s);
-       q = pushLng(mb, q, block_size);
-       q = pushArgument(mb, q, var_next_block);
-       int var_nread = getDestVar(q);
-
-       q = newStmt(mb, "calc", ">");
-       q->barrier = LEAVEsymbol;
-       setReturnArgument(q, var_read_barrier);
-       q = pushArgument(mb, q, var_nread);
-       q = pushLng(mb, q, 0);
-
-       q = newStmt(mb, "streams", "close");
-       q = pushArgument(mb, q, var_s);
-
-       q = newAssignment(mb);
-       setReturnArgument(q, var_s);
-       q = pushNil(mb, q, streams_type);
-
-       // END READ BLOCK
-       q = newAssignment(mb);
-       q->barrier = EXITsymbol;
-       getDestVar(q) = var_read_barrier;
-
-       emit_send(mb, var_stream_channel, streams_type, var_s);
-
-       int var_our_block = emit_receive(mb, var_block_channel, bte_bat_type);
-       int var_our_skip_amount = emit_receive(mb, var_skip_amounts_channel, 
TYPE_int);
-
-       q = newStmt(mb, "aggr", "count");
-       q = pushArgument(mb, q, var_our_block);
-       int var_our_count = getDestVar(q);
-
-       q = newStmt(mb, "aggr", "count");
-       q = pushArgument(mb, q, var_next_block);
-       int var_next_count = getDestVar(q);
-
-       q = newStmt(mb, "calc", "+");
-       q = pushArgument(mb, q, var_our_count);
-       q = pushArgument(mb, q, var_next_count);
-       int var_total_count = getDestVar(q);
-
-       q = newStmt(mb, "calc", "==");
-       q->barrier = LEAVEsymbol;
-       setReturnArgument(q, var_loop_barrier);
-       q = pushArgument(mb, q, var_total_count);
-       q = pushLng(mb, q, 0);
-
-       q = newStmt(mb, "copy", "fixlines");
-       q = pushReturn(mb, q, newTmpVariable(mb, TYPE_int));
-       q = pushArgument(mb, q, var_our_block);
-       q = pushArgument(mb, q, var_our_skip_amount);
-       q = pushArgument(mb, q, var_next_block);
-       q = pushArgument(mb, q, var_line_sep);
-       q = pushArgument(mb, q, var_quote_char);
-       q = pushArgument(mb, q, var_escape);
-       int var_our_line_count = getArg(q, 0);
-       int var_next_skip_amount = getArg(q, 1);
-
-       emit_send(mb, var_block_channel, bte_bat_type, var_next_block);
-       emit_send(mb, var_skip_amounts_channel, TYPE_int, var_next_skip_amount);
-
-       int var_claim_token = emit_receive(mb, var_claim_channel, TYPE_bit);
-
-       q = newStmt(mb, "sql", "claim");
-       q = pushReturn(mb, q, newTmpVariable(mb, newBatType(TYPE_oid)));
-       q = pushArgument(mb, q, be->mvc_var);
-       q = pushStr(mb, q, schema_name);
-       q = pushStr(mb, q, table_name);
-       q = pushArgument(mb, q, var_our_line_count);
-       int var_position = getArg(q, 0);
-       int var_positions = getArg(q, 1);
-
-       emit_send(mb, var_claim_channel, TYPE_bit, var_claim_token);
-
-       //
-       q = newStmt(mb, "calc", "==");
-       q->barrier = REDOsymbol;
-       getDestVar(q) = var_loop_barrier;
-       q = pushArgument(mb, q, var_our_line_count);
-       q = pushLng(mb, q, 0);
-
-       assert(column_count > 0);
-       q = newStmt(mb, "copy", "splitlines");
-       setDestType(mb, q, int_bat_type);
-       for (int i = 1; i < column_count; i++) {
-               int v = newTmpVariable(mb, int_bat_type);
-               q = pushReturn(mb, q, v);
-       }
-       q = pushArgument(mb, q, var_our_block);
-       q = pushArgument(mb, q, var_our_skip_amount);
-       q = pushArgument(mb, q, var_our_line_count);
-       q = pushArgument(mb, q, var_col_sep);
-       q = pushArgument(mb, q, var_line_sep);
-       q = pushArgument(mb, q, var_quote_char);
-       q = pushArgument(mb, q, var_null_representation);
-       q = pushArgument(mb, q, var_escape);
-       InstrPtr splitlines_instr = q;
-
-       int i = 0;
-       for (node *n = table->columns->l->h; n != NULL; n = n->next) {
-               int var_indices = getArg(splitlines_instr, i++);
-
-               sql_column *col = n->data;
-               sql_type *type = col->type.type;
-               const char *column_name = col->base.name;
-
-               switch (type->eclass) {
-                       case EC_DEC:
-                               q = newStmt(mb, "copy", "parse_decimal");
-                               q = pushArgument(mb, q, var_our_block);
-                               q = pushArgument(mb, q, var_indices);
-                               q = pushInt(mb, q, col->type.digits);
-                               q = pushInt(mb, q, col->type.scale);
-                               break;
-                       default:
-                               q = newStmt(mb, "copy", "parse_generic");
-                               q = pushArgument(mb, q, var_our_block);
-                               q = pushArgument(mb, q, var_indices);
-                               q = pushNil(mb, q, col->type.type->localtype);
-                               break;
-               }
-               int var_converted = getDestVar(q);
-
-               q = newStmt(mb, "sql", "append");
-               q = pushArgument(mb, q, be->mvc_var);
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to