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