Changeset: ec95511a4033 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=ec95511a4033
Modified Files:
        sql/include/sql_catalog.h
        sql/server/rel_distribute.c
        sql/server/rel_dump.c
        sql/server/rel_optimizer.c
        sql/server/rel_schema.c
        sql/server/rel_select.c
        sql/server/rel_select.h
        sql/server/sql_parser.y
        sql/server/sql_scan.c
        sql/storage/store.c
Branch: default
Log Message:

add support for 'replica' tables


diffs (truncated from 381 to 300 lines):

diff --git a/sql/include/sql_catalog.h b/sql/include/sql_catalog.h
--- a/sql/include/sql_catalog.h
+++ b/sql/include/sql_catalog.h
@@ -117,7 +117,8 @@ typedef enum temp_t {
        SQL_DECLARED_TABLE,     /* variable inside a stored procedure */
        SQL_MERGE_TABLE,
        SQL_STREAM,
-       SQL_REMOTE
+       SQL_REMOTE,
+       SQL_REPLICA_TABLE
 } temp_t;
 
 typedef enum commit_action_t { 
@@ -378,15 +379,17 @@ typedef enum table_types {
        tt_generated = 2,       /* generated (functions can be sql or c-code) */
        tt_merge_table = 3,     /* multiple tables form one table */
        tt_stream = 4,          /* stream */
-       tt_remote = 5           /* stored on a remote server */
+       tt_remote = 5,          /* stored on a remote server */
+       tt_replica_table = 6    /* multiple replica of the same table */
 } table_types;
 
-#define isTable(x)     (x->type==tt_table)
-#define isView(x)      (x->type==tt_view)
-#define isGenerated(x)  (x->type==tt_generated)
-#define isMergeTable(x) (x->type==tt_merge_table)
-#define isStream(x)    (x->type==tt_stream)
-#define isRemote(x)    (x->type==tt_remote)
+#define isTable(x)       (x->type==tt_table)
+#define isView(x)        (x->type==tt_view)
+#define isGenerated(x)    (x->type==tt_generated)
+#define isMergeTable(x)   (x->type==tt_merge_table)
+#define isStream(x)      (x->type==tt_stream)
+#define isRemote(x)      (x->type==tt_remote)
+#define isReplicaTable(x) (x->type==tt_replica_table)
 
 typedef struct sql_table {
        sql_base base;
diff --git a/sql/server/rel_distribute.c b/sql/server/rel_distribute.c
--- a/sql/server/rel_distribute.c
+++ b/sql/server/rel_distribute.c
@@ -30,6 +30,67 @@
 #include "sql_env.h"
 
 static sql_rel *
+replica(mvc *sql, sql_rel *rel, char *uri) 
+{
+       if (!rel)
+               return rel;
+
+       switch (rel->op) {
+       case op_basetable: {
+               sql_table *t = rel->l;
+
+               if (isReplicaTable(t)) {
+                       node *n;
+
+                       /* replace by the replica which matches the uri */
+                       for (n = t->tables.set->h; n; n = n->next) {
+                               sql_table *p = n->data;
+
+                               if (isRemote(p) && strcmp(uri, p->query) == 0) {
+                                       rel->l = p;
+                                       break;
+                               }
+                       }
+               }
+       }
+       case op_table:
+               break;
+       case op_join: 
+       case op_left: 
+       case op_right: 
+       case op_full: 
+
+       case op_semi: 
+       case op_anti: 
+
+       case op_union: 
+       case op_inter: 
+       case op_except: 
+               rel->l = replica(sql, rel->l, uri);
+               rel->r = replica(sql, rel->r, uri);
+               break;
+       case op_project:
+       case op_select: 
+       case op_groupby: 
+       case op_topn: 
+       case op_sample: 
+               rel->l = replica(sql, rel->l, uri);
+               break;
+       case op_ddl: 
+               rel->l = replica(sql, rel->l, uri);
+               if (rel->r)
+                       rel->r = replica(sql, rel->r, uri);
+               break;
+       case op_insert:
+       case op_update:
+       case op_delete:
+               rel->r = replica(sql, rel->r, uri);
+               break;
+       }
+       return rel;
+}
+
+static sql_rel *
 distribute(mvc *sql, sql_rel *rel) 
 {
        sql_rel *l = NULL, *r = NULL;
@@ -67,6 +128,13 @@ distribute(mvc *sql, sql_rel *rel)
                r = rel->r = distribute(sql, rel->r);
 
                if (l && (pl = find_prop(l->p, PROP_REMOTE)) != NULL &&
+                          r && (pr = find_prop(r->p, PROP_REMOTE)) == NULL) {
+                       r = rel->r = distribute(sql, replica(sql, rel->r, 
pl->value));
+               } else if (l && (pl = find_prop(l->p, PROP_REMOTE)) == NULL &&
+                          r && (pr = find_prop(r->p, PROP_REMOTE)) != NULL) {
+                       l = rel->l = distribute(sql, replica(sql, rel->l, 
pr->value));
+               }
+               if (l && (pl = find_prop(l->p, PROP_REMOTE)) != NULL &&
                    r && (pr = find_prop(r->p, PROP_REMOTE)) != NULL && 
                    strcmp(pl->value, pr->value) == 0) {
                        l->p = prop_remove(l->p, pl);
diff --git a/sql/server/rel_dump.c b/sql/server/rel_dump.c
--- a/sql/server/rel_dump.c
+++ b/sql/server/rel_dump.c
@@ -100,7 +100,9 @@ exp_print(mvc *sql, stream *fout, sql_ex
                        if (atom_type(a)->type->localtype == TYPE_ptr) {
                                sql_table *t = a->data.val.pval;
                                mnstr_printf(fout, "%s(%s)", 
-                                       
isStream(t)?"stream":isMergeTable(t)?"merge table":"table",
+                                       isStream(t)?"stream":
+                                       isMergeTable(t)?"merge table":
+                                       isReplicaTable(t)?"replica 
table":"table",
                                        t->base.name);
                        } else {
                                char *t = sql_subtype_string(atom_type(a));
diff --git a/sql/server/rel_optimizer.c b/sql/server/rel_optimizer.c
--- a/sql/server/rel_optimizer.c
+++ b/sql/server/rel_optimizer.c
@@ -4154,7 +4154,7 @@ rel_remove_unused(mvc *sql, sql_rel *rel
        case op_basetable: {
                sql_table *t = rel->l;
 
-               if (isMergeTable(t)) 
+               if (isMergeTable(t) || isReplicaTable(t)) 
                        return rel;
        }
        case op_table:
@@ -5136,7 +5136,7 @@ rel_merge_table_rewrite(int *changes, mv
                        if (t->tables.set) {
                                for (n = t->tables.set->h; n; n = n->next) {
                                        sql_table *pt = n->data;
-                                       sql_rel *prel = _rel_basetable(sql->sa, 
pt, tname);
+                                       sql_rel *prel = rel_basetable(sql, pt, 
tname);
                                        node *n, *m;
 
                                        /* rename (mostly the idxs) */
diff --git a/sql/server/rel_schema.c b/sql/server/rel_schema.c
--- a/sql/server/rel_schema.c
+++ b/sql/server/rel_schema.c
@@ -162,7 +162,8 @@ mvc_create_table_as_subquery( mvc *sql, 
        char *n;
        int tt =(temp == SQL_REMOTE)?tt_remote:
                (temp == SQL_STREAM)?tt_stream:
-                ((temp == SQL_MERGE_TABLE)?tt_merge_table:tt_table);
+               (temp == SQL_MERGE_TABLE)?tt_merge_table:
+               (temp == SQL_REPLICA_TABLE)?tt_replica_table:tt_table;
 
        sql_table *t = mvc_create_table(sql, s, tname, tt, 0, 
SQL_DECLARED_TABLE, commit_action, -1);
        if ((n = as_subquery( sql, t, sq, column_spec)) != NULL) {
@@ -600,7 +601,7 @@ table_element(mvc *sql, symbol *s, sql_s
 {
        int res = SQL_OK;
 
-       if (alter && (isView(t) || (isMergeTable(t) && s->token != SQL_TABLE && 
s->token != SQL_DROP_TABLE) || (isTable(t) && (s->token == SQL_TABLE || 
s->token == SQL_DROP_TABLE)) )){
+       if (alter && (isView(t) || ((isMergeTable(t) || isReplicaTable(t)) && 
s->token != SQL_TABLE && s->token != SQL_DROP_TABLE) || (isTable(t) && 
(s->token == SQL_TABLE || s->token == SQL_DROP_TABLE)) )){
                char *msg = "";
 
                switch (s->token) {
@@ -634,7 +635,8 @@ table_element(mvc *sql, symbol *s, sql_s
                }
                sql_error(sql, 02, "ALTER TABLE: cannot %s %s '%s'\n",
                                msg, 
-                               isMergeTable(t)?"MERGE TABLE":"VIEW",
+                               isMergeTable(t)?"MERGE TABLE":
+                               isReplicaTable(t)?"REPLICA TABLE":"VIEW",
                                t->base.name);
                return SQL_ERR;
        }
@@ -793,7 +795,8 @@ rel_create_table(mvc *sql, sql_schema *s
        int create = (!instantiate && !deps);
        int tt = (temp == SQL_REMOTE)?tt_remote:
                 (temp == SQL_STREAM)?tt_stream:
-                ((temp == SQL_MERGE_TABLE)?tt_merge_table:tt_table);
+                (temp == SQL_MERGE_TABLE)?tt_merge_table:
+                (temp == SQL_REPLICA_TABLE)?tt_replica_table:tt_table;
 
        (void)create;
        if (sname && !(s = mvc_bind_schema(sql, sname)))
diff --git a/sql/server/rel_select.c b/sql/server/rel_select.c
--- a/sql/server/rel_select.c
+++ b/sql/server/rel_select.c
@@ -432,9 +432,10 @@ rel_copy( sql_allocator *sa, sql_rel *i 
 }
 
 sql_rel *
-_rel_basetable(sql_allocator *sa, sql_table *t, char *atname)
+rel_basetable(mvc *sql, sql_table *t, char *atname)
 {
        node *cn;
+       sql_allocator *sa = sql->sa;
        sql_rel *rel = rel_create(sa);
        char *tname = t->base.name;
 
@@ -473,40 +474,6 @@ _rel_basetable(sql_allocator *sa, sql_ta
 }
 
 sql_rel *
-rel_basetable(mvc *sql, sql_table *t, char *tname)
-{
-#if 0
-       if (isMergeTable(t)) {
-               /* instantiate merge tabel */
-               sql_rel *rel = NULL;
-
-               if (sql->emode == m_deps) {
-                       rel = _rel_basetable(sql->sa, t, tname);
-               } else {
-                       node *n;
-
-                       if (list_empty(t->tables.set)) 
-                               rel = _rel_basetable(sql->sa, t, tname);
-                       if (t->tables.set) {
-                               for (n = t->tables.set->h; n; n = n->next) {
-                                       sql_table *pt = n->data;
-                                       sql_rel *prel = rel_basetable(sql, pt, 
tname);
-                                       if (rel) { 
-                                               rel = rel_setop(sql->sa, rel, 
prel, op_union);
-                                               rel->exps = 
rel_projections(sql, rel, NULL, 1, 1);
-                                       } else {
-                                               rel = prel;
-                                       }
-                               }
-                       }
-               }
-               return rel;
-       }
-#endif
-       return _rel_basetable(sql->sa, t, tname);
-}
-
-sql_rel *
 rel_table_func(sql_allocator *sa, sql_rel *l, sql_exp *f, list *exps)
 {
        sql_rel *rel = rel_create(sa);
diff --git a/sql/server/rel_select.h b/sql/server/rel_select.h
--- a/sql/server/rel_select.h
+++ b/sql/server/rel_select.h
@@ -45,7 +45,6 @@ extern void rel_select_add_exp(sql_rel *
 extern sql_rel *rel_select(sql_allocator *sa, sql_rel *l, sql_exp *e);
 extern sql_rel *rel_select_copy(sql_allocator *sa, sql_rel *l, list *exps);
 extern sql_rel *rel_basetable(mvc *sql, sql_table *t, char *tname);
-extern sql_rel *_rel_basetable(sql_allocator *sa, sql_table *t, char *atname);
 extern sql_rel *rel_recursive_func(sql_allocator *sa, list *exps);
 extern sql_rel *rel_table_func(sql_allocator *sa, sql_rel *l, sql_exp *f, list 
*exps);
 extern sql_rel *rel_relational_func(sql_allocator *sa, sql_rel *l, list *exps);
diff --git a/sql/server/sql_parser.y b/sql/server/sql_parser.y
--- a/sql/server/sql_parser.y
+++ b/sql/server/sql_parser.y
@@ -539,7 +539,7 @@ CONTINUE CURRENT CURSOR FOUND GOTO GO LA
 SQLCODE SQLERROR UNDER WHENEVER
 */
 
-%token TEMPORARY STREAM MERGE REMOTE
+%token TEMPORARY STREAM MERGE REMOTE REPLICA
 %token<sval> ASC DESC AUTHORIZATION
 %token CHECK CONSTRAINT CREATE
 %token TYPE PROCEDURE FUNCTION AGGREGATE RETURNS EXTERNAL sqlNAME DECLARE
@@ -1333,6 +1333,16 @@ table_def:
          append_int(l, commit_action);
          append_string(l, NULL);
          $$ = _symbol_create_list( SQL_CREATE_TABLE, l ); }
+ |  REPLICA TABLE qname table_content_source 
+       { int commit_action = CA_COMMIT, tpe = SQL_REPLICA_TABLE;
+         dlist *l = L();
+
+         append_int(l, tpe);
+         append_list(l, $3);
+         append_symbol(l, $4);
+         append_int(l, commit_action);
+         append_string(l, NULL);
+         $$ = _symbol_create_list( SQL_CREATE_TABLE, l ); }
  /* mapi:monetdb://host:port/database (assumed monetdb/monetdb) */
  |  REMOTE TABLE qname table_content_source ON STRING
        { int commit_action = CA_COMMIT, tpe = SQL_REMOTE;
diff --git a/sql/server/sql_scan.c b/sql/server/sql_scan.c
--- a/sql/server/sql_scan.c
+++ b/sql/server/sql_scan.c
@@ -198,6 +198,7 @@ scanner_init_keywords(void)
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list

Reply via email to