From 132b7ecd1282e7829662a4cdda0040647eb81900 Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Mon, 11 Jan 2016 16:06:55 -0800 Subject: [PATCH 01/10] Factor out parallel reading; start to set up semi-parallel reading --- geojson.c | 227 ++++++++++++++++++++++++++++++++++++++---------------- 1 file changed, 159 insertions(+), 68 deletions(-) diff --git a/geojson.c b/geojson.c index 35f06c0d..19ef9fa5 100644 --- a/geojson.c +++ b/geojson.c @@ -902,6 +902,84 @@ void *run_sort(void *v) { return NULL; } +void do_read_parallel(char *map, long long len, long long initial_offset, 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) { + long long segs[CPUS + 1]; + segs[0] = 0; + segs[CPUS] = len; + + int i; + for (i = 1; i < CPUS; i++) { + segs[i] = len * i / CPUS; + + while (segs[i] < len && map[segs[i]] != '\n') { + segs[i]++; + } + } + + long long layer_seq[CPUS]; + for (i = 0; i < CPUS; i++) { + // To preserve feature ordering, unique id for each segment + // begins with that segment's offset into the input + layer_seq[i] = segs[i] + initial_offset; + } + + struct parse_json_args pja[CPUS]; + pthread_t pthreads[CPUS]; + + for (i = 0; i < CPUS; i++) { + pja[i].jp = json_begin_map(map + segs[i], segs[i + 1] - segs[i]); + pja[i].reading = reading; + pja[i].layer_seq = &layer_seq[i]; + pja[i].progress_seq = progress_seq; + pja[i].metapos = &reader[i].metapos; + pja[i].geompos = &reader[i].geompos; + pja[i].indexpos = &reader[i].indexpos; + pja[i].exclude = exclude; + pja[i].include = include; + pja[i].exclude_all = exclude_all; + pja[i].metafile = reader[i].metafile; + pja[i].geomfile = reader[i].geomfile; + pja[i].indexfile = reader[i].indexfile; + pja[i].poolfile = reader[i].poolfile; + pja[i].treefile = reader[i].treefile; + pja[i].fname = fname; + pja[i].maxzoom = maxzoom; + pja[i].basezoom = basezoom; + pja[i].layer = source < nlayers ? source : 0; + pja[i].droprate = droprate; + pja[i].file_bbox = reader[i].file_bbox; + pja[i].segment = i; + pja[i].initialized = &initialized[i]; + pja[i].initial_x = &initial_x[i]; + pja[i].initial_y = &initial_y[i]; + + if (pthread_create(&pthreads[i], NULL, run_parse_json, &pja[i]) != 0) { + perror("pthread_create"); + exit(EXIT_FAILURE); + } + } + + for (i = 0; i < CPUS; i++) { + void *retval; + + if (pthread_join(pthreads[i], &retval) != 0) { + perror("pthread_join"); + } + + free(pja[i].jp->source); + json_end(pja[i].jp); + } +} + +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) { + 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) { int ret = EXIT_SUCCESS; @@ -1048,71 +1126,10 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } if (map != NULL && map != MAP_FAILED) { - long long segs[CPUS + 1]; - segs[0] = 0; - segs[CPUS] = st.st_size - off; + do_read_parallel(map, st.st_size - off, 0, reading, reader, &progress_seq, exclude, include, exclude_all, fname, maxzoom, basezoom, source, nlayers, droprate, initialized, initial_x, initial_y); - int i; - for (i = 1; i < CPUS; i++) { - segs[i] = off + (st.st_size - off) * i / CPUS; - - while (segs[i] < st.st_size && map[segs[i]] != '\n') { - segs[i]++; - } - } - - long long layer_seq[CPUS]; - for (i = 0; i < CPUS; i++) { - // To preserve feature ordering, unique id for each segment - // begins with that segment's offset into the input - layer_seq[i] = segs[i]; - } - - struct parse_json_args pja[CPUS]; - pthread_t pthreads[CPUS]; - - for (i = 0; i < CPUS; i++) { - pja[i].jp = json_begin_map(map + segs[i], segs[i + 1] - segs[i]); - pja[i].reading = reading; - pja[i].layer_seq = &layer_seq[i]; - pja[i].progress_seq = &progress_seq; - pja[i].metapos = &reader[i].metapos; - pja[i].geompos = &reader[i].geompos; - pja[i].indexpos = &reader[i].indexpos; - pja[i].exclude = exclude; - pja[i].include = include; - pja[i].exclude_all = exclude_all; - pja[i].metafile = reader[i].metafile; - pja[i].geomfile = reader[i].geomfile; - pja[i].indexfile = reader[i].indexfile; - pja[i].poolfile = reader[i].poolfile; - pja[i].treefile = reader[i].treefile; - pja[i].fname = fname; - pja[i].maxzoom = maxzoom; - pja[i].basezoom = basezoom; - pja[i].layer = source < nlayers ? source : 0; - pja[i].droprate = droprate; - pja[i].file_bbox = reader[i].file_bbox; - pja[i].segment = i; - pja[i].initialized = &initialized[i]; - pja[i].initial_x = &initial_x[i]; - pja[i].initial_y = &initial_y[i]; - - if (pthread_create(&pthreads[i], NULL, run_parse_json, &pja[i]) != 0) { - perror("pthread_create"); - exit(EXIT_FAILURE); - } - } - - for (i = 0; i < CPUS; i++) { - void *retval; - - if (pthread_join(pthreads[i], &retval) != 0) { - perror("pthread_join"); - } - - free(pja[i].jp->source); - json_end(pja[i].jp); + if (munmap(map, st.st_size - off) != 0) { + perror("munmap source file"); } } else { FILE *fp = fdopen(fd, "r"); @@ -1122,10 +1139,84 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max continue; } - long long layer_seq = 0; - json_pull *jp = json_begin_file(fp); - parse_json(jp, reading, &layer_seq, &progress_seq, &reader[0].metapos, &reader[0].geompos, &reader[0].indexpos, exclude, include, exclude_all, reader[0].metafile, reader[0].geomfile, reader[0].indexfile, reader[0].poolfile, reader[0].treefile, fname, maxzoom, basezoom, source < nlayers ? source : 0, droprate, reader[0].file_bbox, 0, &initialized[0], &initial_x[0], &initial_y[0]); - json_end(jp); + if (read_parallel) { + // Serial reading of chunks that are then parsed in parallel + + char readname[strlen(tmpdir) + strlen("/read.XXXXXXXX") + 1]; + sprintf(readname, "%s%s", tmpdir, "/read.XXXXXXXX"); + int readfd = mkstemp(readname); + if (readfd < 0) { + perror(readname); + exit(EXIT_FAILURE); + } + FILE *readfp = fdopen(readfd, "w"); + if (readfp == NULL) { + perror(readname); + exit(EXIT_FAILURE); + } + unlink(readname); + + int c; + volatile int is_parsing = 0; + long long ahead = 0; + long long initial_offset = 0; + pthread_t previous_reader; + + while ((c = getc(fp)) != EOF) { + putc(c, readfp); + ahead++; + + if (c == '\n' && ahead > 100000 && is_parsing == 0) { + if (initial_offset != 0) { + if (pthread_join(previous_reader, NULL) != 0) { + perror("pthread_join"); + exit(EXIT_FAILURE); + } + } + + fclose(readfp); + start_parsing(readfd, initial_offset, &is_parsing, &previous_reader); + + initial_offset += ahead; + ahead = 0; + + sprintf(readname, "%s%s", tmpdir, "/read.XXXXXXXX"); + int readfd = mkstemp(readname); + if (readfd < 0) { + perror(readname); + exit(EXIT_FAILURE); + } + FILE *readfp = fdopen(readfd, "w"); + if (readfp == NULL) { + perror(readname); + exit(EXIT_FAILURE); + } + unlink(readname); + } + } + + if (initial_offset != 0) { + if (pthread_join(previous_reader, NULL) != 0) { + perror("pthread_join"); + exit(EXIT_FAILURE); + } + } + + fclose(readfp); + start_parsing(readfd, initial_offset, &is_parsing, &previous_reader); + + if (pthread_join(previous_reader, NULL) != 0) { + perror("pthread_join"); + } + } else { + // Plain serial reading + + long long layer_seq = 0; + json_pull *jp = json_begin_file(fp); + parse_json(jp, reading, &layer_seq, &progress_seq, &reader[0].metapos, &reader[0].geompos, &reader[0].indexpos, exclude, include, exclude_all, reader[0].metafile, reader[0].geomfile, reader[0].indexfile, reader[0].poolfile, reader[0].treefile, fname, maxzoom, basezoom, source < nlayers ? source : 0, droprate, reader[0].file_bbox, 0, &initialized[0], &initial_x[0], &initial_y[0]); + json_end(jp); + } + fclose(fp); } } From 2d1657794517b1321609a7f2ded9f375aebfe0ac Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Mon, 11 Jan 2016 16:52:45 -0800 Subject: [PATCH 02/10] Starts but crashes --- geojson.c | 91 +++++++++++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 82 insertions(+), 9 deletions(-) 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 From 333956ce42e60c66c80873ff165b8b4c1d6dc964 Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Mon, 11 Jan 2016 17:29:06 -0800 Subject: [PATCH 03/10] Fix crashes --- geojson.c | 36 +++++++++++++++++++++++++----------- 1 file changed, 25 insertions(+), 11 deletions(-) diff --git a/geojson.c b/geojson.c index c7596dee..f515f929 100644 --- a/geojson.c +++ b/geojson.c @@ -973,6 +973,7 @@ void do_read_parallel(char *map, long long len, long long initial_offset, const struct start_parsing_arg { int fd; + FILE *fp; long long offset; long long len; volatile int *is_parsing; @@ -997,6 +998,15 @@ struct start_parsing_arg { void *run_start_parsing(void *v) { struct start_parsing_arg *a = v; + struct stat st; + if (fstat(a->fd, &st) != 0) { + perror("stat read temp"); + } + if (a->len != st.st_size) { + printf("%lld vs %lld\n", a->len, st.st_size); + } + a->len = st.st_size; + char *map = mmap(NULL, a->len, PROT_READ, MAP_PRIVATE, a->fd, 0); if (map == NULL || map == MAP_FAILED) { perror("map intermediate input"); @@ -1008,14 +1018,17 @@ void *run_start_parsing(void *v) { if (munmap(map, a->len) != 0) { perror("munmap source file"); } + if (fclose(a->fp) != 0) { + perror("close source file"); + } - a->is_parsing = 0; + *(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) { +void start_parsing(int fd, FILE *fp, 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. @@ -1024,6 +1037,7 @@ void start_parsing(int fd, long long offset, long long len, volatile int *is_par struct start_parsing_arg *spa = malloc(sizeof(struct start_parsing_arg)); spa->fd = fd; + spa->fp = fp; spa->offset = offset; spa->len = len; spa->is_parsing = is_parsing; @@ -1190,7 +1204,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max if (fstat(fd, &st) == 0) { off = lseek(fd, 0, SEEK_CUR); if (off >= 0) { - map = mmap(NULL, st.st_size - off, PROT_READ, MAP_PRIVATE, fd, off); + // map = mmap(NULL, st.st_size - off, PROT_READ, MAP_PRIVATE, fd, off); } } } @@ -1219,7 +1233,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max perror(readname); exit(EXIT_FAILURE); } - FILE *readfp = fopen(readname, "w"); + FILE *readfp = fdopen(readfd, "w"); if (readfp == NULL) { perror(readname); exit(EXIT_FAILURE); @@ -1236,7 +1250,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max putc(c, readfp); ahead++; - if (c == '\n' && ahead > 100000 && is_parsing == 0) { + if (c == '\n' && ahead > 1000000 && is_parsing == 0) { if (initial_offset != 0) { if (pthread_join(previous_reader, NULL) != 0) { perror("pthread_join"); @@ -1244,19 +1258,19 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } } - fclose(readfp); - 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); + fflush(readfp); + start_parsing(readfd, readfp, 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; sprintf(readname, "%s%s", tmpdir, "/read.XXXXXXXX"); - int readfd = mkstemp(readname); + readfd = mkstemp(readname); if (readfd < 0) { perror(readname); exit(EXIT_FAILURE); } - FILE *readfp = fdopen(readfd, "w"); + readfp = fdopen(readfd, "w"); if (readfp == NULL) { perror(readname); exit(EXIT_FAILURE); @@ -1272,10 +1286,10 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } } - fclose(readfp); + fflush(readfp); 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); + start_parsing(readfd, readfp, 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"); From 9d6ece5bbc1450efd09eb6d8678340f044d16aa0 Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 11:47:46 -0800 Subject: [PATCH 04/10] Buffered reading makes it faster than the single-threaded version --- geojson.c | 31 +++++++++++++++++++------------ 1 file changed, 19 insertions(+), 12 deletions(-) diff --git a/geojson.c b/geojson.c index f515f929..06d70c50 100644 --- a/geojson.c +++ b/geojson.c @@ -741,11 +741,6 @@ void parse_json(json_pull *jp, const char *reading, long long *layer_seq, volati /* XXX check for any non-features in the outer object */ } - - if (!quiet) { - fprintf(stderr, " \r"); - // (stderr, "Read 10000.00 million features\r", *progress_seq / 1000000.0); - } } struct parse_json_args { @@ -1003,7 +998,7 @@ void *run_start_parsing(void *v) { perror("stat read temp"); } if (a->len != st.st_size) { - printf("%lld vs %lld\n", a->len, st.st_size); + fprintf(stderr, "wrong number of bytes in teomporary: %lld vs %lld\n", a->len, st.st_size); } a->len = st.st_size; @@ -1204,7 +1199,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max if (fstat(fd, &st) == 0) { off = lseek(fd, 0, SEEK_CUR); if (off >= 0) { - // map = mmap(NULL, st.st_size - off, PROT_READ, MAP_PRIVATE, fd, off); + map = mmap(NULL, st.st_size - off, PROT_READ, MAP_PRIVATE, fd, off); } } } @@ -1240,17 +1235,21 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } unlink(readname); - int c; volatile int is_parsing = 0; long long ahead = 0; long long initial_offset = 0; pthread_t previous_reader; - while ((c = getc(fp)) != EOF) { - putc(c, readfp); - ahead++; +#define READ_BUF 2000 +#define PARSE_MIN 10000000 + char buf[READ_BUF]; + int n; - if (c == '\n' && ahead > 1000000 && is_parsing == 0) { + while ((n = fread(buf, sizeof(char), READ_BUF, fp)) > 0) { + fwrite(buf, sizeof(char), n, readfp); + ahead += n; + + if (buf[n - 1] == '\n' && ahead > PARSE_MIN && is_parsing == 0) { if (initial_offset != 0) { if (pthread_join(previous_reader, NULL) != 0) { perror("pthread_join"); @@ -1278,6 +1277,9 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max unlink(readname); } } + if (n < 0) { + perror(reading); + } if (initial_offset != 0) { if (pthread_join(previous_reader, NULL) != 0) { @@ -1308,6 +1310,11 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } } + if (!quiet) { + fprintf(stderr, " \r"); + // (stderr, "Read 10000.00 million features\r", *progress_seq / 1000000.0); + } + for (i = 0; i < CPUS; i++) { fclose(reader[i].metafile); fclose(reader[i].geomfile); From 0680236e46084e919c172dc935e246c8378cf990 Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 12:18:05 -0800 Subject: [PATCH 05/10] Fix warning --- geojson.c | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/geojson.c b/geojson.c index 06d70c50..b966e3ba 100644 --- a/geojson.c +++ b/geojson.c @@ -998,7 +998,7 @@ void *run_start_parsing(void *v) { perror("stat read temp"); } if (a->len != st.st_size) { - fprintf(stderr, "wrong number of bytes in teomporary: %lld vs %lld\n", a->len, st.st_size); + fprintf(stderr, "wrong number of bytes in temporary: %lld vs %lld\n", a->len, (long long) st.st_size); } a->len = st.st_size; @@ -1246,7 +1246,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max int n; while ((n = fread(buf, sizeof(char), READ_BUF, fp)) > 0) { - fwrite(buf, sizeof(char), n, readfp); + fwrite_check(buf, sizeof(char), n, readfp, reading); ahead += n; if (buf[n - 1] == '\n' && ahead > PARSE_MIN && is_parsing == 0) { From e4afaa7a27968f4361b4081535f7d63d80d36b0f Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 12:31:17 -0800 Subject: [PATCH 06/10] Renaming in the hope of clarity --- geojson.c | 64 +++++++++++++++++++++++++++---------------------------- 1 file changed, 32 insertions(+), 32 deletions(-) diff --git a/geojson.c b/geojson.c index b966e3ba..9936ef37 100644 --- a/geojson.c +++ b/geojson.c @@ -966,7 +966,7 @@ void do_read_parallel(char *map, long long len, long long initial_offset, const } } -struct start_parsing_arg { +struct read_parallel_arg { int fd; FILE *fp; long long offset; @@ -990,8 +990,8 @@ struct start_parsing_arg { unsigned *initial_y; }; -void *run_start_parsing(void *v) { - struct start_parsing_arg *a = v; +void *run_read_parallel(void *v) { + struct read_parallel_arg *a = v; struct stat st; if (fstat(a->fd, &st) != 0) { @@ -1023,37 +1023,37 @@ void *run_start_parsing(void *v) { return NULL; } -void start_parsing(int fd, FILE *fp, 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) { +void start_parsing(int fd, FILE *fp, long long offset, long long len, volatile int *is_parsing, pthread_t *parallel_parser, 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->fp = fp; - spa->offset = offset; - spa->len = len; - spa->is_parsing = is_parsing; + struct read_parallel_arg *rpa = malloc(sizeof(struct read_parallel_arg)); + rpa->fd = fd; + rpa->fp = fp; + rpa->offset = offset; + rpa->len = len; + rpa->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; + rpa->reading = reading; + rpa->reader = reader; + rpa->progress_seq = progress_seq; + rpa->exclude = exclude; + rpa->include = include; + rpa->exclude_all = exclude_all; + rpa->fname = fname; + rpa->maxzoom = maxzoom; + rpa->basezoom = basezoom; + rpa->source = source; + rpa->nlayers = nlayers; + rpa->droprate = droprate; + rpa->initialized = initialized; + rpa->initial_x = initial_x; + rpa->initial_y = initial_y; - if (pthread_create(previous_reader, NULL, run_start_parsing, spa) != 0) { + if (pthread_create(parallel_parser, NULL, run_read_parallel, rpa) != 0) { perror("pthread_create"); exit(EXIT_FAILURE); } @@ -1238,7 +1238,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max volatile int is_parsing = 0; long long ahead = 0; long long initial_offset = 0; - pthread_t previous_reader; + pthread_t parallel_parser; #define READ_BUF 2000 #define PARSE_MIN 10000000 @@ -1251,14 +1251,14 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max if (buf[n - 1] == '\n' && ahead > PARSE_MIN && is_parsing == 0) { if (initial_offset != 0) { - if (pthread_join(previous_reader, NULL) != 0) { + if (pthread_join(parallel_parser, NULL) != 0) { perror("pthread_join"); exit(EXIT_FAILURE); } } fflush(readfp); - start_parsing(readfd, readfp, 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); + start_parsing(readfd, readfp, initial_offset, ahead, &is_parsing, ¶llel_parser, reading, reader, &progress_seq, exclude, include, exclude_all, fname, maxzoom, basezoom, source, nlayers, droprate, initialized, initial_x, initial_y); initial_offset += ahead; ahead = 0; @@ -1282,7 +1282,7 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max } if (initial_offset != 0) { - if (pthread_join(previous_reader, NULL) != 0) { + if (pthread_join(parallel_parser, NULL) != 0) { perror("pthread_join"); exit(EXIT_FAILURE); } @@ -1291,9 +1291,9 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max fflush(readfp); if (ahead > 0) { - start_parsing(readfd, readfp, 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); + start_parsing(readfd, readfp, initial_offset, ahead, &is_parsing, ¶llel_parser, 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) { + if (pthread_join(parallel_parser, NULL) != 0) { perror("pthread_join"); } } From ca97c5ec6d1d22a05f60bd62ce85be24e12aa6c0 Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 12:36:12 -0800 Subject: [PATCH 07/10] Update docs --- README.md | 2 +- man/tippecanoe.1 | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 62140e10..da500cf4 100644 --- a/README.md +++ b/README.md @@ -68,7 +68,7 @@ Options * -P: Use multiple threads to read different parts of each input file at once. This will only work if the input is line-delimited JSON with each Feature on its own line, because it knows nothing of the top-level structure around the Features. - In addition, it only works if the input is a named file that can be mapped into memory + Performance will be better if the input is a named file that can be mapped into memory rather than a stream that can only be read sequentially. ### Zoom levels and resolution diff --git a/man/tippecanoe.1 b/man/tippecanoe.1 index 765a7408..15c0a9e2 100644 --- a/man/tippecanoe.1 +++ b/man/tippecanoe.1 @@ -74,7 +74,7 @@ specified, the files are all merged into the single named layer. \-P: Use multiple threads to read different parts of each input file at once. This will only work if the input is line\-delimited JSON with each Feature on its own line, because it knows nothing of the top\-level structure around the Features. -In addition, it only works if the input is a named file that can be mapped into memory +Performance will be better if the input is a named file that can be mapped into memory rather than a stream that can only be read sequentially. .RE .SS Zoom levels and resolution From ecae14e2d4d5e7789f5a7ef6c78a585654f07daa Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 14:12:56 -0800 Subject: [PATCH 08/10] Stabilize feature order between the different reading methods --- geojson.c | 16 +++++++++------- 1 file changed, 9 insertions(+), 7 deletions(-) diff --git a/geojson.c b/geojson.c index 9936ef37..2e8e4b69 100644 --- a/geojson.c +++ b/geojson.c @@ -274,7 +274,8 @@ struct index { long long start; long long end; unsigned long long index; - int segment; + short segment; + unsigned long long seq : (64 - 16); // pack with segment to stay in 32 bytes }; int indexcmp(const void *v1, const void *v2) { @@ -287,6 +288,12 @@ int indexcmp(const void *v1, const void *v2) { return 1; } + if (i1->seq < i2->seq) { + return -1; + } else if (i1->seq > i2->seq) { + return 1; + } + return 0; } @@ -599,6 +606,7 @@ int serialize_geometry(json_object *geometry, json_object *properties, const cha index.start = geomstart; index.end = *geompos; index.segment = segment; + index.seq = *layer_seq; // Calculate the center even if off the edge of the plane, // and then mask to bring it back into the addressable area @@ -1430,12 +1438,6 @@ int read_json(int argc, char **argv, char *fname, const char *layername, int max fprintf(stderr, "Sorting %lld features\n", (long long) indexpos / bytes); } - // XXX On machines with different page sizes, doing the sorting - // in different-sized chunks can cause features with the same - // index (i.e., the same bbox) to appear in different orders - // because the sort is unstable. This doesn't seem worth spending - // more memory to fix, but could make tests unstable. - int page = sysconf(_SC_PAGESIZE); long long unit = (50 * 1024 * 1024 / bytes) * bytes; while (unit % page != 0) { From 83322d8e35a29afc72c2005d5fc0addc65cd94c4 Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 14:16:17 -0800 Subject: [PATCH 09/10] Guard against unlikely overflow --- geojson.c | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/geojson.c b/geojson.c index 2e8e4b69..b8cb9ec5 100644 --- a/geojson.c +++ b/geojson.c @@ -69,6 +69,11 @@ void init_cpus() { CPUS = 1; } + // Guard against short struct index.segment + if (CPUS > 32767) { + CPUS = 32767; + } + // Round down to a power of 2 CPUS = 1 << (int) (log(CPUS) / log(2)); From 872df4bd9f0d25e64823799d8db1f972d388edab Mon Sep 17 00:00:00 2001 From: Eric Fischer Date: Tue, 12 Jan 2016 14:27:05 -0800 Subject: [PATCH 10/10] Bump version number --- CHANGELOG.md | 5 +++++ geojson.c | 2 +- version.h | 2 +- 3 files changed, 7 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e5ff3cc6..304ebcc4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,8 @@ +## 1.7.0 + +* Parallel processing of input with -P works with streamed input too +* Error handling if unsupported options given to -p or -a + ## 1.6.4 * Fix crashing bug when layers are being merged with -l diff --git a/geojson.c b/geojson.c index b8cb9ec5..09e74dc6 100644 --- a/geojson.c +++ b/geojson.c @@ -280,7 +280,7 @@ struct index { long long end; unsigned long long index; short segment; - unsigned long long seq : (64 - 16); // pack with segment to stay in 32 bytes + unsigned long long seq : (64 - 16); // pack with segment to stay in 32 bytes }; int indexcmp(const void *v1, const void *v2) { diff --git a/version.h b/version.h index 14cb23d8..2c7c8d74 100644 --- a/version.h +++ b/version.h @@ -1 +1 @@ -#define VERSION "tippecanoe v1.6.4\n" +#define VERSION "tippecanoe v1.7.0\n"