Changeset: b59803cefee9 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/b59803cefee9
Modified Files:
        sql/backends/monet5/rel_copy.c
Branch: copyparpipe
Log Message:

Extract ON SERVER specific code into function


diffs (truncated from 322 to 300 lines):

diff --git a/sql/backends/monet5/rel_copy.c b/sql/backends/monet5/rel_copy.c
--- a/sql/backends/monet5/rel_copy.c
+++ b/sql/backends/monet5/rel_copy.c
@@ -47,6 +47,122 @@ emit_send(MalBlkPtr mb, int var_channel,
        q = pushNil(mb, q, tpe);
 }
 
+struct loop_vars {
+       int loop_barrier;
+       int our_block;
+       int our_skip_amount;
+       int our_line_count;
+};
+
+
+static void
+emit_onserver(
+       MalBlkPtr mb, struct loop_vars *loop_vars,
+       int var_fname, int block_size, int margin,
+       int var_line_sep, int var_quote_char, int var_escape)
+{
+       InstrPtr q;
+
+       int streams_type = ATOMindex("streams");
+       int bte_bat_type = newBatType(TYPE_bte);
+
+       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);
+
+
+       // START LOOP
+       q = newAssignment(mb);
+       q->barrier = BARRIERsymbol;
+       q = pushBit(mb, q, true);
+       loop_vars->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);
+
+       loop_vars->our_block = emit_receive(mb, var_block_channel, 
bte_bat_type);
+       loop_vars->our_skip_amount = emit_receive(mb, var_skip_amounts_channel, 
TYPE_int);
+
+       q = newStmt(mb, "aggr", "count");
+       q = pushArgument(mb, q, loop_vars->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, loop_vars->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, loop_vars->our_block);
+       q = pushArgument(mb, q, loop_vars->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);
+       loop_vars->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);
+}
+
+
 stmt *
 rel2bin_copyparpipe(backend *be, sql_rel *rel, list *refs, sql_exp *copyfrom)
 {
@@ -55,14 +171,13 @@ rel2bin_copyparpipe(backend *be, sql_rel
        const int block_size = 1024 * 1024;
        const int margin = 8 * 1024;
 
+       struct loop_vars loop_vars;
        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
@@ -108,104 +223,12 @@ rel2bin_copyparpipe(backend *be, sql_rel
        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);
+       emit_onserver(mb, &loop_vars, var_fname, block_size, margin, 
var_line_sep, var_quote_char, var_escape);
 
        int var_claim_token = emit_receive(mb, var_claim_channel, TYPE_bit);
 
@@ -214,17 +237,16 @@ rel2bin_copyparpipe(backend *be, sql_rel
        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);
+       q = pushArgument(mb, q, loop_vars.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);
+       getDestVar(q) = loop_vars.loop_barrier;
+       q = pushArgument(mb, q, loop_vars.our_line_count);
        q = pushLng(mb, q, 0);
 
        assert(column_count > 0);
@@ -234,9 +256,9 @@ rel2bin_copyparpipe(backend *be, sql_rel
                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, loop_vars.our_block);
+       q = pushArgument(mb, q, loop_vars.our_skip_amount);
+       q = pushArgument(mb, q, loop_vars.our_line_count);
        q = pushArgument(mb, q, var_col_sep);
        q = pushArgument(mb, q, var_line_sep);
        q = pushArgument(mb, q, var_quote_char);
@@ -255,14 +277,14 @@ rel2bin_copyparpipe(backend *be, sql_rel
                switch (type->eclass) {
                        case EC_DEC:
                                q = newStmt(mb, "copy", "parse_decimal");
-                               q = pushArgument(mb, q, var_our_block);
+                               q = pushArgument(mb, q, loop_vars.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, loop_vars.our_block);
                                q = pushArgument(mb, q, var_indices);
                                q = pushNil(mb, q, col->type.type->localtype);
                                break;
@@ -282,19 +304,19 @@ rel2bin_copyparpipe(backend *be, sql_rel
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to