diff --git a/geojson.c b/geojson.c index 19ef9fa5..c7596dee 100644 --- a/geojson.c +++ b/geojson.c @@ -971,13 +971,83 @@ void do_read_parallel(char *map, long long len, long long initial_offset, const } } -void start_parsing(int fd, long long offset, volatile int *is_parsing, pthread_t *previous_reader) { -#if 0 - if (pthread_create(previous_reader, NULL, run_start_parsing, start_parsing_args) == 0) { +struct start_parsing_arg { + int fd; + long long offset; + long long len; + volatile int *is_parsing; + + const char *reading; + struct reader *reader; + volatile long long *progress_seq; + struct pool *exclude; + struct pool *include; + int exclude_all; + char *fname; + int maxzoom; + int basezoom; + int source; + int nlayers; + double droprate; + int *initialized; + unsigned *initial_x; + unsigned *initial_y; +}; + +void *run_start_parsing(void *v) { + struct start_parsing_arg *a = v; + + char *map = mmap(NULL, a->len, PROT_READ, MAP_PRIVATE, a->fd, 0); + if (map == NULL || map == MAP_FAILED) { + perror("map intermediate input"); + exit(EXIT_FAILURE); + } + + do_read_parallel(map, a->len, a->offset, a->reading, a->reader, a->progress_seq, a->exclude, a->include, a->exclude_all, a->fname, a->maxzoom, a->basezoom, a->source, a->nlayers, a->droprate, a->initialized, a->initial_x, a->initial_y); + + if (munmap(map, a->len) != 0) { + perror("munmap source file"); + } + + a->is_parsing = 0; + free(a); + + return NULL; +} + +void start_parsing(int fd, long long offset, long long len, volatile int *is_parsing, pthread_t *previous_reader, const char *reading, struct reader *reader, volatile long long *progress_seq, struct pool *exclude, struct pool *include, int exclude_all, char *fname, int maxzoom, int basezoom, int source, int nlayers, double droprate, int *initialized, unsigned *initial_x, unsigned *initial_y) { + // This has to kick off an intermediate thread to start the parser threads, + // so the main thread can get back to reading the next input stage while + // the intermediate thread waits for the completion of the parser threads. + + *is_parsing = 1; + + struct start_parsing_arg *spa = malloc(sizeof(struct start_parsing_arg)); + spa->fd = fd; + spa->offset = offset; + spa->len = len; + spa->is_parsing = is_parsing; + + spa->reading = reading; + spa->reader = reader; + spa->progress_seq = progress_seq; + spa->exclude = exclude; + spa->include = include; + spa->exclude_all = exclude_all; + spa->fname = fname; + spa->maxzoom = maxzoom; + spa->basezoom = basezoom; + spa->source = source; + spa->nlayers = nlayers; + spa->droprate = droprate; + spa->initialized = initialized; + spa->initial_x = initial_x; + spa->initial_y = initial_y; + + if (pthread_create(previous_reader, NULL, run_start_parsing, spa) != 0) { perror("pthread_create"); exit(EXIT_FAILURE); } -#endif } int read_json(int argc, char **argv, char *fname, const char *layername, int maxzoom, int minzoom, int basezoom, double basezoom_marker_width, sqlite3 *outdb, struct pool *exclude, struct pool *include, int exclude_all, double droprate, int buffer, const char *tmpdir, double gamma, char *prevent, char *additional, int read_parallel) { @@ -1149,7 +1219,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max perror(readname); exit(EXIT_FAILURE); } - FILE *readfp = fdopen(readfd, "w"); + FILE *readfp = fopen(readname, "w"); if (readfp == NULL) { perror(readname); exit(EXIT_FAILURE); @@ -1175,7 +1245,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } fclose(readfp); - start_parsing(readfd, initial_offset, &is_parsing, &previous_reader); + start_parsing(readfd, initial_offset, ahead, &is_parsing, &previous_reader, reading, reader, &progress_seq, exclude, include, exclude_all, fname, maxzoom, basezoom, source, nlayers, droprate, initialized, initial_x, initial_y); initial_offset += ahead; ahead = 0; @@ -1203,10 +1273,13 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } fclose(readfp); - start_parsing(readfd, initial_offset, &is_parsing, &previous_reader); - if (pthread_join(previous_reader, NULL) != 0) { - perror("pthread_join"); + if (ahead > 0) { + start_parsing(readfd, initial_offset, ahead, &is_parsing, &previous_reader, reading, reader, &progress_seq, exclude, include, exclude_all, fname, maxzoom, basezoom, source, nlayers, droprate, initialized, initial_x, initial_y); + + if (pthread_join(previous_reader, NULL) != 0) { + perror("pthread_join"); + } } } else { // Plain serial reading