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]