Changeset: 3339240300f3 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=3339240300f3
Modified Files:
sql/backends/monet5/miniseed/registrar.c
Branch: DVframework
Log Message:
Kept the sortedness of file ids for multi-threaded registration.
diffs (truncated from 312 to 300 lines):
diff --git a/sql/backends/monet5/miniseed/registrar.c
b/sql/backends/monet5/miniseed/registrar.c
--- a/sql/backends/monet5/miniseed/registrar.c
+++ b/sql/backends/monet5/miniseed/registrar.c
@@ -44,9 +44,9 @@ typedef struct {
lng loop_start;
lng loop_end;
int mode; /* carries to the thread */
- Client cntxt; /* carries to the thread */
- int *function_created;
+ mvc *mvc;
int base_file_id;
+ temp_container* tc;
} thread_argv;
lng get_line_num(str filename);
@@ -67,14 +67,12 @@ str mseed_register_and_mount_segments_mo
int concatenate_strs(str* words_to_concat, int num_words_to_concat, str*
ret_concatenated);
str prepare_insertion(Client cntxt, temp_container* tc);
str insert_into_vault(Client cntxt, temp_container* tc);
+str clean_up_after_insertion(Client cntxt, temp_container* tc);
str SQLstatementIntern(Client c, str *expr, str nme, int execute, bit output);
str register_clean_up(temp_container* tc);
void *register_files(void *args);
// static str runSQLQuery(Client c, char **result, char * query);
-pthread_mutex_t create_lock;
-pthread_mutex_t insert_lock;
-
double timetol = -1.0; /* Time tolerance for continuous traces */
double sampratetol = -1.0; /* Sample rate tolerance for continuous traces */
int verbose = 1;
@@ -881,7 +879,7 @@ str prepare_insertion(Client cntxt, temp
}
q = (str)GDKmalloc(512*sizeof(char));
- sprintf(q, "CREATE FUNCTION %s_%s_reg(ticket bigint, table_idx
int) RETURNS table(%s) external name registrar.register_table;\n",
tc->schema_name, tc->table_names[t], concatenated);
+ sprintf(q, "CREATE FUNCTION %s.%s_%s_reg(ticket bigint,
table_idx int) RETURNS table(%s) external name registrar.register_table;\n",
tc->schema_name, tc->schema_name, tc->table_names[t], concatenated);
if((msg
=SQLstatementIntern(cntxt,&q,"registrar.create.function",TRUE,FALSE))!=
MAL_SUCCEED)
{/* create function query not succeeded, what to do */
@@ -913,26 +911,47 @@ str insert_into_vault(Client cntxt, temp
for(t = 0; t < tc->num_tables; t++)
{
str q = (str)GDKmalloc(512*sizeof(char));
- str s = (str)GDKmalloc(512*sizeof(char));
- sprintf(q, "INSERT INTO %s.%s SELECT * FROM %s_%s_reg("LLFMT",
%d);\n", tc->schema_name, tc->table_names[t], tc->schema_name,
tc->table_names[t], ticket, t);
+ sprintf(q, "INSERT INTO %s.%s SELECT * FROM
%s.%s_%s_reg("LLFMT", %d);\n", tc->schema_name, tc->table_names[t],
tc->schema_name, tc->schema_name, tc->table_names[t], ticket, t);
if((msg
=SQLstatementIntern(cntxt,&q,"registrar.insert",TRUE,FALSE))!= MAL_SUCCEED)
{/* insert into query not succeeded, what to do */
return msg;
}
GDKfree(q);
+ }
+
+ return MAL_SUCCEED;
+}
+
+/*
+ * forms and executes sql 'DROP FUNCTION' queries for each of
+ * the tables_to_be_filled which are in temp_container tc. Drops the new sql
+ * functions created by a prepare_insertion call.
+ *
+ * returns error or MAL_SUCCEED.
+ */
+str clean_up_after_insertion(Client cntxt, temp_container* tc)
+{
+ /* form a sql query str like this: */
+ /* INSERT INTO mseed.files SELECT * FROM mseed_files_reg(ticket,
table_idx); */
+
+ int t;
+ str msg;
+
+ for(t = 0; t < tc->num_tables; t++)
+ {
+ str s = (str)GDKmalloc(512*sizeof(char));
+ sprintf(s, "DROP FUNCTION %s.%s_%s_reg(BIGINT, INTEGER);\n",
tc->schema_name, tc->schema_name, tc->table_names[t]);
- sprintf(s, "DROP FUNCTION %s_%s_reg(BIGINT, INTEGER);\n",
tc->schema_name, tc->table_names[t]);
-
- if((msg
=SQLstatementIntern(cntxt,&s,"registrar.insert",TRUE,FALSE))!= MAL_SUCCEED)
+ if((msg
=SQLstatementIntern(cntxt,&s,"registrar.clean_up_after_insertion",TRUE,FALSE))!=
MAL_SUCCEED)
{/* drop function not succeeded, what to do */
return msg;
}
GDKfree(s);
}
-
+
return MAL_SUCCEED;
}
@@ -1314,7 +1333,7 @@ str mseed_register_segments_mode(int ass
str ch = (str) GDKmalloc(2*sizeof(char));
ch[1] = '\0';
- printf(""LLFMT".file\n", file_counter);
+// printf(""LLFMT".file\n", file_counter);
file_counter++;
mstl = mstl_init (NULL);
@@ -1435,7 +1454,7 @@ str mseed_register_segments_mode(int ass
gap = -(((double)(seg->endtime -
seg->starttime)/HPTMODULUS) + delta);
}
- printf("%d. segment, frq: %lf, sampcnt: %lld\n",
segment_id_fake, seg->samprate, (long long int)seg->samplecnt);
+// printf("%d. segment, frq: %lf, sampcnt: %lld\n",
segment_id_fake, seg->samprate, (long long int)seg->samplecnt);
if ((aBAT =
BATdescriptor(ret_tc->tables_columns[1].column_bats[4])) == NULL)
throw(MAL, "mseed_register",
RUNTIME_OBJECT_MISSING);
@@ -1787,7 +1806,7 @@ str mseed_register_and_mount_segments_mo
void *register_files(void *args)
{
- temp_container *tc;
+// temp_container *tc;
lng i;
str err = NULL;
int start, finish;
@@ -1796,12 +1815,14 @@ void *register_files(void *args)
thread_argv targv = *((thread_argv*)args);
/* create temp_container */
- tc = (temp_container*)GDKmalloc(sizeof(temp_container));
- assert(tc != NULL);
- if(targv.mode == 0)
- err = mseed_create_temp_container_segments_mode(tc); /*
depending on design can get different argument(s) */
- else
- err =
mseed_create_temp_container_with_data_tables_segments_mode(tc); /* depending on
design can get different argument(s) */
+// tc = (temp_container*)GDKmalloc(sizeof(temp_container));
+// assert(tc != NULL);
+
+ err = create_temp_container_segments_mode(targv.mvc, targv.tc,
targv.mode); /* depending on design can get different argument(s) */
+// if(targv.mode == 0)
+// err = mseed_create_temp_container_segments_mode(targv.tc); /*
depending on design can get different argument(s) */
+// else
+// err =
mseed_create_temp_container_with_data_tables_segments_mode(targv.tc); /*
depending on design can get different argument(s) */
if(err != MAL_SUCCEED)
{/* temp_container creation failed, what to do */
throw(MAL,"registrar.register_repo", "temp_container creation
failed in thread %d: %s\n", targv.tid, err);
@@ -1816,7 +1837,7 @@ void *register_files(void *args)
{
for(i = targv.loop_start; i < targv.loop_end; i++)
{
- err = mseed_register_segments_mode(current_file_id,
targv.file_paths[i], tc);
+ err = mseed_register_segments_mode(current_file_id,
targv.file_paths[i], targv.tc);
if(err != MAL_SUCCEED)
{/* current file cannot be registered, what to do */
/*throw(MAL,"registrar.register_repo", "Current
file cannot be registered: %s\n", err); */
@@ -1830,7 +1851,7 @@ void *register_files(void *args)
{
for(i = targv.loop_start; i < targv.loop_end; i++)
{
- err =
mseed_register_and_mount_segments_mode(current_file_id, targv.file_paths[i],
tc);
+ err =
mseed_register_and_mount_segments_mode(current_file_id, targv.file_paths[i],
targv.tc);
if(err != MAL_SUCCEED)
{/* current file cannot be registered, what to do */
/* throw(MAL,"registrar.register_repo",
"Current file cannot be registered: %s\n", err); */
@@ -1843,37 +1864,6 @@ void *register_files(void *args)
finish = GDKms();
printf("# In thread %d, time for extraction and transformation of
(meta-)data: %d milliseconds\n", targv.tid, finish - start);
- pthread_mutex_lock(&create_lock);
- if(*targv.function_created == 0)
- {
- /* prepare sql functions for inserting temp_container into
tables_to_be_filled */
- err = prepare_insertion(targv.cntxt, tc);
- if(err != MAL_SUCCEED)
- {/* preparing the insertion failed, what to do */
- throw(MAL,"registrar.register_repo", "Insertion prepare
failed in thread %d: %s\n", targv.tid, err);
- }
- *targv.function_created = 1;
- }
- pthread_mutex_unlock(&create_lock);
-
- pthread_mutex_lock(&insert_lock);
- start = GDKms();
- /* insert temp_container into tables_to_be_filled */
- err = insert_into_vault(targv.cntxt, tc);
- if(err != MAL_SUCCEED)
- {/* inserting the temp_container into one of the tables failed, what to
do */
- throw(MAL,"registrar.register_repo", "Inserting the
temp_container into one of the tables failed in thread %d: %s\n", targv.tid,
err);
- }
- finish = GDKms();
- printf("# In thread %d, time for loading of (meta-)data: %d
milliseconds\n", targv.tid, finish - start);
- pthread_mutex_unlock(&insert_lock);
-
- err = register_clean_up(tc);
- if(err != MAL_SUCCEED)
- {/* inserting the temp_container into one of the tables failed, what to
do */
- throw(MAL,"registrar.register_repo", "Cleaning up the temp_container
failed in thread %d: %s\n", targv.tid, err);
- }
-
return NULL;
}
@@ -1896,10 +1886,10 @@ str register_repo(Client cntxt, MalBlkPt
str *file_paths = NULL;
lng num_file_paths;
temp_container *tc;
+ temp_container **thread_tcs;
lng i;
str err = NULL;
int start, finish, total_start, total_finish;
- int function_created = 0;
mvc *m = NULL;
int max_file_id = -1;
int current_file_id;
@@ -1957,15 +1947,16 @@ str register_repo(Client cntxt, MalBlkPt
{
/* multi-threaded */
+ int j; /* thread index */
lng loop_start = 0;
lng num_file_paths_per_thread = num_file_paths / num_threads;
pthread_t *threads =
(pthread_t*)GDKmalloc(num_threads*sizeof(pthread_t));
thread_argv *targvs =
(thread_argv*)GDKmalloc(num_threads*sizeof(thread_argv));
- int j;
-
- if (pthread_mutex_init(&insert_lock, NULL) != 0 ||
pthread_mutex_init(&create_lock, NULL) != 0)
+ thread_tcs =
(temp_container**)GDKmalloc(num_threads*sizeof(temp_container*));
+
+ for(j = 0; j < num_threads; j++)
{
- throw(MAL,"registrar.register_repo", "mutex init
failed\n");
+ thread_tcs[j] =
(temp_container*)GDKmalloc(num_threads*sizeof(temp_container));
}
for(j = 0; j < num_threads; j++)
@@ -1980,8 +1971,9 @@ str register_repo(Client cntxt, MalBlkPt
targvs[j].base_file_id = max_file_id + 1 +
targvs[j].loop_start;
loop_start = targvs[j].loop_end;
targvs[j].mode = mode;
- targvs[j].cntxt = cntxt;
- targvs[j].function_created = &function_created;
+ targvs[j].mvc = m;
+
+ targvs[j].tc = thread_tcs[j];
pthread_create(&threads[j], NULL, register_files, (void
*) &targvs[j]);
}
@@ -1991,8 +1983,42 @@ str register_repo(Client cntxt, MalBlkPt
pthread_join(threads[j], NULL);
}
- pthread_mutex_destroy(&create_lock);
- pthread_mutex_destroy(&insert_lock);
+ /* prepare sql functions for inserting temp_container into
tables_to_be_filled */
+ err = prepare_insertion(cntxt, thread_tcs[0]);
+ if(err != MAL_SUCCEED)
+ {/* preparing the insertion failed, what to do */
+ throw(MAL,"registrar.register_repo", "Insertion prepare
failed: %s\n", err);
+ }
+
+ start = GDKms();
+ for(j = 0; j < num_threads; j++)
+ {
+ /* insert temp_container into tables_to_be_filled */
+ err = insert_into_vault(cntxt, thread_tcs[j]);
+ if(err != MAL_SUCCEED)
+ {/* inserting the temp_container into one of the tables
failed, what to do */
+ throw(MAL,"registrar.register_repo", "Inserting
the temp_container into one of the tables failed: %s\n", err);
+ }
+ }
+
+ err = clean_up_after_insertion(cntxt, thread_tcs[0]);
+ if(err != MAL_SUCCEED)
+ {/* cleaning up the insertion failed, what to do */
+ throw(MAL,"registrar.register_repo", "Insertion clean
up failed: %s\n", err);
+ }
+
+ for(j = 0; j < num_threads; j++)
+ {
+ err = register_clean_up(thread_tcs[j]);
+ if(err != MAL_SUCCEED)
+ {/* inserting the temp_container into one of the tables
failed, what to do */
+ throw(MAL,"registrar.register_repo", "Cleaning
up the temp_container failed: %s\n", err);
+ }
+ }
+
+ finish = GDKms();
+ printf("# Time for loading of (meta-)data: %d milliseconds\n",
finish - start);
+
GDKfree(targvs);
GDKfree(threads);
@@ -2058,14 +2084,21 @@ str register_repo(Client cntxt, MalBlkPt
{/* inserting the temp_container into one of the tables failed,
what to do */
throw(MAL,"registrar.register_repo", "Inserting the
temp_container into one of the tables failed: %s\n", err);
}
- finish = GDKms();
- printf("# Time for loading of (meta-)data: %d milliseconds\n",
finish - start);
+
+ err = clean_up_after_insertion(cntxt, tc);
+ if(err != MAL_SUCCEED)
+ {/* cleaning up the insertion failed, what to do */
+ throw(MAL,"registrar.register_repo", "Insertion clean
up failed: %s\n", err);
+ }
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list