Changeset: 6237bba18938 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/6237bba18938
Branch: default
Log Message:
Merge with group-commit.
diffs (truncated from 672 to 300 lines):
diff --git a/clients/Tests/exports.stable.out b/clients/Tests/exports.stable.out
--- a/clients/Tests/exports.stable.out
+++ b/clients/Tests/exports.stable.out
@@ -543,9 +543,9 @@ gdk_return log_bat_transient(logger *lg,
gdk_return log_constant(logger *lg, int type, ptr val, log_id id, lng offset,
lng cnt);
gdk_return log_delta(logger *lg, BAT *uid, BAT *uval, log_id id);
gdk_return log_sequence(logger *lg, int seq, lng id);
-gdk_return log_tdone(logger *lg, ulng commit_ts);
gdk_return log_tend(logger *lg);
-gdk_return log_tstart(logger *lg, bool flush);
+gdk_return log_tflush(logger *lg, ulng log_file_id, ulng commit_ts);
+gdk_return log_tstart(logger *lg, bool flushnow, ulng *log_file_id);
gdk_return logger_activate(logger *lg);
lng logger_changes(logger *lg);
logger *logger_create(int debug, const char *fn, const char *logdir, int
version, preversionfix_fptr prefuncp, postversionfix_fptr postfuncp, void
*funcdata);
diff --git a/gdk/gdk_logger.c b/gdk/gdk_logger.c
--- a/gdk/gdk_logger.c
+++ b/gdk/gdk_logger.c
@@ -2048,7 +2048,12 @@ logger_load(int debug, const char *fn, c
logbat_destroy(lg->seqs_id);
logbat_destroy(lg->seqs_val);
logbat_destroy(lg->dseqs);
+ ATOMIC_DESTROY(&lg->refcount);
MT_lock_destroy(&lg->lock);
+ MT_lock_destroy(&lg->rotation_lock);
+ MT_sema_destroy(&lg->flush_queue_semaphore);
+ MT_lock_destroy(&lg->flush_lock);
+ MT_lock_destroy(&lg->flush_queue_lock);
GDKfree(lg->fn);
GDKfree(lg->dir);
GDKfree(lg->local_dir);
@@ -2108,11 +2113,21 @@ logger_new(int debug, const char *fn, co
GDKfree(lg);
return NULL;
}
- MT_lock_init(&lg->lock, fn);
if (lg->debug & 1) {
fprintf(stderr, "#logger_new dir set to %s\n", lg->dir);
}
+ ATOMIC_INIT(&lg->refcount, 0);
+ MT_lock_init(&lg->lock, fn);
+ MT_lock_init(&lg->rotation_lock, "rotation_lock");
+ MT_sema_init(&lg->flush_queue_semaphore, FLUSH_QUEUE_SIZE,
"flush_queue_semaphore");
+ MT_lock_init(&lg->flush_lock, "flush_lock");
+ MT_lock_init(&lg->flush_queue_lock, "flush_queue_lock");
+
+ // flush variables
+ lg->flush_queue_begin = 0;
+ lg->flush_queue_length = 0;
+
if (logger_load(debug, fn, logdir, lg, filename) == GDK_SUCCEED) {
return lg;
}
@@ -2156,7 +2171,12 @@ logger_destroy(logger *lg)
logbat_destroy(lg->catalog_lid);
logger_unlock(lg);
}
+ ATOMIC_DESTROY(&lg->refcount);
MT_lock_destroy(&lg->lock);
+ MT_lock_destroy(&lg->rotation_lock);
+ MT_sema_destroy(&lg->flush_queue_semaphore);
+ MT_lock_destroy(&lg->flush_lock);
+ MT_lock_destroy(&lg->flush_queue_lock);
GDKfree(lg->fn);
GDKfree(lg->dir);
GDKfree(lg->buf);
@@ -2215,6 +2235,7 @@ logger_cleanup_range(logger *lg)
gdk_return
logger_activate(logger *lg)
{
+ MT_lock_set(&lg->rotation_lock);
logger_lock(lg);
if (lg->end > 0 && lg->saved_id+1 == lg->id) {
lg->id++;
@@ -2222,10 +2243,12 @@ logger_activate(logger *lg)
/* start new file */
if (logger_open_output(lg) != GDK_SUCCEED) {
logger_unlock(lg);
+ MT_lock_unset(&lg->rotation_lock);
return GDK_FAIL;
}
}
logger_unlock(lg);
+ MT_lock_unset(&lg->rotation_lock);
return GDK_SUCCEED;
}
@@ -2361,6 +2384,7 @@ log_constant(logger *lg, int type, ptr v
(!is_row && !mnstr_writeLng(lg->output_log, nr)) ||
(!is_row && mnstr_write(lg->output_log, &tpe, 1, 1) != 1) ||
(!is_row && !mnstr_writeLng(lg->output_log, offset))) {
+ (void) ATOMIC_DEC(&lg->refcount);
ok = GDK_FAIL;
goto bailout;
}
@@ -2507,6 +2531,7 @@ internal_log_bat(logger *lg, BAT *b, log
bailout:
if (ok != GDK_SUCCEED) {
+ (void) ATOMIC_DEC(&lg->refcount);
const char *err = mnstr_peek_error(lg->output_log);
TRC_CRITICAL(GDK, "write failed%s%s\n", err ? ": " : "", err ?
err : "");
}
@@ -2527,6 +2552,8 @@ log_bat_persists(logger *lg, BAT *b, log
if (logger_add_bat(lg, b, id, -1) != GDK_SUCCEED) {
logger_unlock(lg);
+ if (!LOG_DISABLED(lg))
+ (void) ATOMIC_DEC(&lg->refcount);
return GDK_FAIL;
}
@@ -2536,6 +2563,7 @@ log_bat_persists(logger *lg, BAT *b, log
if (log_write_format(lg, &l) != GDK_SUCCEED ||
mnstr_write(lg->output_log, &ta, 1, 1) != 1) {
logger_unlock(lg);
+ (void) ATOMIC_DEC(&lg->refcount);
return GDK_FAIL;
}
}
@@ -2544,6 +2572,8 @@ log_bat_persists(logger *lg, BAT *b, log
fprintf(stderr, "#persists id (%d) bat (%d)\n", id,
b->batCacheid);
gdk_return r = internal_log_bat(lg, b, id, 0, BATcount(b), 0);
logger_unlock(lg);
+ if (r != GDK_SUCCEED)
+ (void) ATOMIC_DEC(&lg->refcount);
return r;
}
@@ -2561,6 +2591,7 @@ log_bat_transient(logger *lg, log_id id)
if (log_write_format(lg, &l) != GDK_SUCCEED) {
TRC_CRITICAL(GDK, "write failed\n");
logger_unlock(lg);
+ (void) ATOMIC_DEC(&lg->refcount);
return GDK_FAIL;
}
}
@@ -2571,6 +2602,8 @@ log_bat_transient(logger *lg, log_id id)
lg->end += BATcount(BBPquickdesc(bid));
gdk_return r = logger_del_bat(lg, bid);
logger_unlock(lg);
+ if (r != GDK_SUCCEED)
+ (void) ATOMIC_DEC(&lg->refcount);
return r;
}
@@ -2596,6 +2629,8 @@ log_delta(logger *lg, BAT *uid, BAT *uva
if (BATtdense(uid)) {
ok = internal_log_bat(lg, uval, id, uid->tseqbase,
BATcount(uval), 1);
logger_unlock(lg);
+ if (!LOG_DISABLED(lg) && ok != GDK_SUCCEED)
+ (void) ATOMIC_DEC(&lg->refcount);
return ok;
}
@@ -2654,6 +2689,7 @@ log_delta(logger *lg, BAT *uid, BAT *uva
if (ok != GDK_SUCCEED) {
const char *err = mnstr_peek_error(lg->output_log);
TRC_CRITICAL(GDK, "write failed%s%s\n", err ? ": " : "", err ?
err : "");
+ (void) ATOMIC_DEC(&lg->refcount);
}
logger_unlock(lg);
return ok;
@@ -2678,7 +2714,11 @@ log_bat_clear(logger *lg, int id)
if (lg->debug & 1)
fprintf(stderr, "#Logged clear %d\n", id);
- return log_write_format(lg, &l);
+
+ gdk_return r = log_write_format(lg, &l);
+ if(r != GDK_SUCCEED)
+ (void) ATOMIC_DEC(&lg->refcount);
+ return r;
}
#define DBLKSZ 8192
@@ -2690,60 +2730,161 @@ log_bat_clear(logger *lg, int id)
static gdk_return
new_logfile(logger *lg)
{
- lng log_large = (GDKdebug & FORCEMITOMASK)?LOG_MINI:LOG_LARGE;
assert(!LOG_DISABLED(lg));
- lng p;
- p = (lng) getfilepos(getFile(lg->output_log));
- if (p == -1)
+
+
+ const lng log_large = (GDKdebug & FORCEMITOMASK)?LOG_MINI:LOG_LARGE;
+
+ gdk_return result = GDK_SUCCEED;
+ MT_lock_set(&lg->rotation_lock);
+ const lng p = (lng) getfilepos(getFile(lg->output_log));
+ if (p == -1) {
+ MT_lock_unset(&lg->rotation_lock);
return GDK_FAIL;
- if (p > log_large || (lg->end*1024) > log_large) {
- lg->id++;
- logger_close_output(lg);
- return logger_open_output(lg);
}
- return GDK_SUCCEED;
+ if (( p > log_large || (lg->end*1024) > log_large )) {
+ if (ATOMIC_GET(&lg->refcount) == 1) {
+ lg->id++;
+ logger_close_output(lg);
+ result = logger_open_output(lg);
+ lg->request_rotation = false;
+ }
+ else {
+ // Delegate wal rotation to next writer or last flusher.
+ lg->request_rotation = true;
+ }
+ }
+ MT_lock_unset(&lg->rotation_lock);
+ return result;
}
gdk_return
log_tend(logger *lg)
{
- logformat l;
-
if (lg->debug & 1)
fprintf(stderr, "#log_tend %d\n", lg->tid);
- l.flag = LOG_END;
- l.id = lg->tid;
- if (lg->flushnow) {
- lg->flushnow = 0;
- return logger_commit(lg);
- }
-
lg->end++;
if (LOG_DISABLED(lg)) {
return GDK_SUCCEED;
}
- if (log_write_format(lg, &l) != GDK_SUCCEED ||
- mnstr_flush(lg->output_log, MNSTR_FLUSH_DATA) ||
- (!(GDKdebug & NOSYNCMASK) && mnstr_fsync(lg->output_log)) ||
- new_logfile(lg) != GDK_SUCCEED) {
- TRC_CRITICAL(GDK, "write failed\n");
- return GDK_FAIL;
+ gdk_return result;
+ logformat l;
+ l.flag = LOG_END;
+ l.id = lg->tid;
+
+ if ((result = log_write_format(lg, &l)) != GDK_SUCCEED)
+ (void) ATOMIC_DEC(&lg->refcount);
+ return result;
+}
+static int
+request_number_flush_queue(logger *lg) {
+ // Semaphore protects ring buffer structure in queue against overflowing
+ static int _number = 0;
+ int result;
+ MT_sema_down(&lg->flush_queue_semaphore);
+ MT_lock_set(&lg->flush_queue_lock);
+ result = ++_number;
+ const int end = (lg->flush_queue_begin + lg->flush_queue_length) %
FLUSH_QUEUE_SIZE;
+ lg->flush_queue[end] = _number;
+ lg->flush_queue_length++;
+ MT_lock_unset(&lg->flush_queue_lock);
+
+ return result;
+}
+
+static void
+left_truncate_flush_queue(logger *lg, int limit) {
+ MT_lock_set(&lg->flush_queue_lock);
+ lg->flush_queue_begin = (lg->flush_queue_begin + limit) %
FLUSH_QUEUE_SIZE;
+ lg->flush_queue_length -= limit;
+ MT_lock_unset(&lg->flush_queue_lock);
+
+ for (int i = 0; i < limit; i++)
+ MT_sema_up(&lg->flush_queue_semaphore);
+}
+
+static int
+number_in_flush_queue(logger *lg, int number) {
+ MT_lock_set(&lg->flush_queue_lock);
+ const int fql = lg->flush_queue_length;
+ MT_lock_unset(&lg->flush_queue_lock);
+ for (int i = 0; i < fql; i++) {
+ const int idx = (lg->flush_queue_begin + i) % FLUSH_QUEUE_SIZE;
+ if (lg->flush_queue[idx] == number) {
+ return 1;
+ }
+ }
+ return 0;
+}
+
+static int
+flush_queue_length(logger *lg) {
+ MT_lock_set(&lg->flush_queue_lock);
+ const int fql = lg->flush_queue_length;
+ MT_lock_unset(&lg->flush_queue_lock);
+ return fql;
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]