Changeset: 23362f63873f for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/23362f63873f
Modified Files:
sql/backends/monet5/sql.c
sql/storage/bat/bat_storage.c
sql/storage/bat/bat_storage.h
sql/storage/store.c
Branch: Jul2021
Log Message:
beter split of mergering updates and garbage collection the delta structures.
Solves crash in concurrent update cases.
diffs (259 lines):
diff --git a/sql/backends/monet5/sql.c b/sql/backends/monet5/sql.c
--- a/sql/backends/monet5/sql.c
+++ b/sql/backends/monet5/sql.c
@@ -2150,6 +2150,13 @@ DELTAsub(bat *result, const bat *col, co
BBPunfix(u->batCacheid);
throw(MAL, "sql.delta", SQLSTATE(HY013)
MAL_MALLOC_FAIL);
}
+ BAT *nres;
+ if ((nres = COLcopy(res, res->ttype, true, TRANSIENT))
== NULL) {
+ BBPunfix(res->batCacheid);
+ throw(MAL, "sql.delta", GDK_EXCEPTION);
+ }
+ BBPunfix(res->batCacheid);
+ res = nres;
ret = BATappend(res, u, cminu, true);
BBPunfix(u->batCacheid);
if (cminu)
diff --git a/sql/storage/bat/bat_storage.c b/sql/storage/bat/bat_storage.c
--- a/sql/storage/bat/bat_storage.c
+++ b/sql/storage/bat/bat_storage.c
@@ -841,13 +841,13 @@ older_delta( sql_delta *d, sql_trans *tr
{
sql_delta *o = d->next;
- while (o) {
- if (o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
+ while (o && !o->cs.merged) {
+ if (o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
break;
else
o = o->next;
}
- if (o && o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
+ if (o && !o->cs.merged && o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
return o;
return NULL;
}
@@ -1699,6 +1699,7 @@ append_col_execute(sql_trans *tr, sql_de
{
int ok = LOG_OK;
+ delta->cs.merged = 0;
if (is_bat) {
BAT *bat = incoming_data;
@@ -2327,7 +2328,7 @@ commit_create_col_( sql_trans *tr, sql_c
delta->cs.ts = commit_ts;
assert(delta->next == NULL);
- if (!delta->cs.alter)
+ if (!delta->cs.alter && !delta->cs.merged)
ok = merge_delta(delta);
delta->cs.alter = 0;
if (!tr->parent)
@@ -2433,7 +2434,8 @@ commit_create_idx_( sql_trans *tr, sql_i
delta->cs.ts = commit_ts;
assert(delta->next == NULL);
- ok = merge_delta(delta);
+ if (!delta->cs.alter && !delta->cs.merged)
+ ok = merge_delta(delta);
if (!tr->parent)
i->base.new = 0;
}
@@ -3139,6 +3141,7 @@ merge_cs( column_storage *cs)
bat_destroy(uv);
}
cs->cleared = 0;
+ cs->merged = 1;
bat_destroy(cur);
return ok;
}
@@ -3146,6 +3149,8 @@ merge_cs( column_storage *cs)
static int
merge_delta( sql_delta *obat)
{
+ if (obat && obat->next && !obat->cs.merged)
+ merge_delta(obat->next);
return merge_cs(&obat->cs);
}
@@ -3217,9 +3222,10 @@ commit_update_col_( sql_trans *tr, sql_c
(void)oldest;
if (isTempTable(c->t)) {
if (commit_ts) { /* commit */
- if (c->t->commit_action == CA_COMMIT ||
c->t->commit_action == CA_PRESERVE)
- ok = merge_delta(delta);
- else /* CA_DELETE as CA_DROP's are gone already (or for
globals are equal to a CA_DELETE) */
+ if (c->t->commit_action == CA_COMMIT ||
c->t->commit_action == CA_PRESERVE) {
+ if (!delta->cs.merged)
+ ok = merge_delta(delta);
+ } else /* CA_DELETE as CA_DROP's are gone already (or
for globals are equal to a CA_DELETE) */
clear_cs(tr, &delta->cs, true,
isTempTable(c->t));
} else { /* rollback */
if (c->t->commit_action == CA_COMMIT/* ||
c->t->commit_action == CA_PRESERVE*/)
@@ -3290,20 +3296,12 @@ commit_update_col( sql_trans *tr, sql_ch
o->next = d->next;
change->cleanup = &tc_gc_rollbacked;
} else if (ok == LOG_OK && !tr->parent) {
- sql_delta *d = delta;
- /* clean up and merge deltas */
- while (delta && delta->cs.ts > oldest) {
+ /* merge deltas */
+ while (delta && delta->cs.ts > oldest)
delta = delta->next;
- }
- if (delta && delta != d) {
- if (delta->next) {
- ok = destroy_delta(delta->next, true);
- delta->next = NULL;
- }
- }
- if (ok == LOG_OK && delta == d && oldest == commit_ts) {
- lock_column(tr->store, c->base.id);
- ok = merge_delta(delta);
+ if (ok == LOG_OK && delta && !delta->cs.merged && delta->cs.ts
<= oldest) {
+ lock_column(tr->store, c->base.id); /* lock for
concurrent updates (appends) */
+ merge_delta(delta);
unlock_column(tr->store, c->base.id);
}
} else if (ok == LOG_OK && tr->parent) /* move delta into older and
cleanup current save points */
@@ -3333,9 +3331,10 @@ commit_update_idx_( sql_trans *tr, sql_i
(void)oldest;
if (isTempTable(i->t)) {
if (commit_ts) { /* commit */
- if (i->t->commit_action == CA_COMMIT ||
i->t->commit_action == CA_PRESERVE)
- ok = merge_delta(delta);
- else /* CA_DELETE as CA_DROP's are gone already */
+ if (i->t->commit_action == CA_COMMIT ||
i->t->commit_action == CA_PRESERVE) {
+ if (!delta->cs.merged)
+ ok = merge_delta(delta);
+ } else /* CA_DELETE as CA_DROP's are gone already */
clear_cs(tr, &delta->cs, true,
isTempTable(i->t));
} else { /* rollback */
if (i->t->commit_action == CA_COMMIT/* ||
i->t->commit_action == CA_PRESERVE*/)
@@ -3375,20 +3374,12 @@ commit_update_idx( sql_trans *tr, sql_ch
o->next = d->next;
change->cleanup = &tc_gc_rollbacked;
} else if (ok == LOG_OK && !tr->parent) {
- sql_delta *d = delta;
- /* clean up and merge deltas */
- while (delta && delta->cs.ts > oldest) {
+ /* merge deltas */
+ while (delta && delta->cs.ts > oldest)
delta = delta->next;
- }
- if (delta && delta != d) {
- if (delta->next) {
- ok = destroy_delta(delta->next, true);
- delta->next = NULL;
- }
- }
- if (ok == LOG_OK && delta == d && oldest == commit_ts) {
- lock_column(tr->store, i->base.id);
- ok = merge_delta(delta);
+ if (ok == LOG_OK && delta && !delta->cs.merged && delta->cs.ts
<= oldest) {
+ lock_column(tr->store, i->base.id); /* lock for
concurrent updates (appends) */
+ merge_delta(delta);
unlock_column(tr->store, i->base.id);
}
} else if (ok == LOG_OK && tr->parent) /* cleanup older save points */
@@ -3500,6 +3491,18 @@ commit_update_del( sql_trans *tr, sql_ch
return ok;
}
+static int
+gc_delta( sql_store Store, sql_change *change, ulng oldest)
+{
+ sqlstore *store = Store;
+ sql_delta *n = change->data;
+ (void)store;
+ (void)oldest;
+
+ destroy_delta(n, true);
+ return 1;
+}
+
/* only rollback (content version) case for now */
static int
gc_col( sqlstore *store, sql_change *change, ulng oldest, bool cleanup)
@@ -3516,11 +3519,24 @@ gc_col( sqlstore *store, sql_change *cha
return 0;
sql_delta *d = (sql_delta*)change->data;
if (d->next) {
+ assert(!cleanup);
if (d->cs.ts > oldest)
return LOG_OK; /* cannot cleanup yet */
- destroy_delta(d->next, true);
+ sql_delta *n = d->next;
+ if (n->cs.ucnt && !n->cs.merged) {
+ lock_column(store, c->base.id); /* lock for concurrent
updates (appends) */
+ merge_delta(n);
+ unlock_column(store, c->base.id);
+ } else if (d && d->cs.ucnt && !d->cs.merged) {
+ lock_column(store, c->base.id); /* lock for concurrent
updates (appends) */
+ merge_delta(d);
+ unlock_column(store, c->base.id);
+ }
d->next = NULL;
+ change->cleanup = &gc_delta;
+ change->data = n;
+ return 0;
}
if (cleanup)
column_destroy(store, c);
@@ -3555,11 +3571,24 @@ gc_idx( sqlstore *store, sql_change *cha
return 0;
sql_delta *d = (sql_delta*)change->data;
if (d->next) {
+ assert(!cleanup);
if (d->cs.ts > oldest)
return LOG_OK; /* cannot cleanup yet */
- destroy_delta(d->next, true);
+ sql_delta *n = d->next;
+ if (n->cs.ucnt && !n->cs.merged) {
+ lock_column(store, i->base.id); /* lock for concurrent
updates (appends) */
+ merge_delta(n);
+ unlock_column(store, i->base.id);
+ } else if (d && d->cs.ucnt && !d->cs.merged) {
+ lock_column(store, i->base.id); /* lock for concurrent
updates (appends) */
+ merge_delta(d);
+ unlock_column(store, i->base.id);
+ }
d->next = NULL;
+ change->cleanup = &gc_delta;
+ change->data = n;
+ return 0;
}
if (cleanup)
idx_destroy(store, i);
diff --git a/sql/storage/bat/bat_storage.h b/sql/storage/bat/bat_storage.h
--- a/sql/storage/bat/bat_storage.h
+++ b/sql/storage/bat/bat_storage.h
@@ -19,6 +19,7 @@ typedef struct column_storage {
int uvbid; /* bat with values of updates */
bool cleared;
bool alter; /* set when the delta is created for an alter
statement */
+ bool merged; /* only merge changes once */
size_t ucnt; /* number of updates */
ulng ts; /* version timestamp */
} column_storage;
diff --git a/sql/storage/store.c b/sql/storage/store.c
--- a/sql/storage/store.c
+++ b/sql/storage/store.c
@@ -5707,11 +5707,10 @@ sql_trans_drop_column(sql_trans *tr, sql
n = nn;
col = next;
} else if (col) { /* if the column to be dropped was found,
decrease the column number for others after it */
- oid rid;
next->colnr--;
if (!isDeclaredTable(t)) {
- rid = store->table_api.column_find_row(tr, cid,
&next->base.id, NULL);
+ oid rid = store->table_api.column_find_row(tr,
cid, &next->base.id, NULL);
assert(!is_oid_nil(rid));
if ((res =
store->table_api.column_update_value(tr, cnr, rid, &next->colnr)))
return res;
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list