Reduce thrashing during feature ingestion and tiling (#56)

* Remove the concept of "separate metadata"

This was an extra level of attribute indirection (features point
to metadata records which point to key and value strings) which was
intended to reduce the size of temporary storage for features with
large numbers of attributes that were also spread across large numbers
of tiles at maxzoom.

For other kinds of features, the extra indirection slowed things down
instead, and, especially when maxzoom guessing was being used, many more
features were having their metadata externalized than could actually
benefit from it.

* Shave a few bytes off temporary files by using more unsigned integers

* Flush stderr after logging progress

* Revert "Shave a few bytes off temporary files by using more unsigned integers"

This reverts commit eef29084ec.

* Limit the size of the string pools and trees to fit in memory

* Add missing #include

* Move the string pool and search tree from mmap to allocated memory

* Sort in allocated rather than mapped memory too

* Also use pread instead of mapping to read in the data to sort

* When the pool gets too big, switch to just the file, not memory

* Switch string pool from memory to disk when memory is 10% full

* Add to-memory versions of the serialization functions

* Crashy work in progress toward compression

* Fix the pointer bug that was causing the crash

* Serialize features into memory rather than straight to disk

* Compress individual features in the temporary files

* Don't need to store the length of the geometry

* Remove per-feature compression; move minzoom back into the object

* Start adding a stream compressor object

* Track file position within fwrite_check()

* Add compressed stream writer functions

* Pull the writing of the serialized feature out to the callers

* Starting toward compression again from a different point

* Hook up more compression functions

* Remove unused code from the other day

* Make enough deflate calls to flush out all the buffered data

* Start on decompression

* Tile number is uncompressed, tile content is compressed

* Work on alternating compressed and uncompressed in decompression

* Closer, but still doesn't work

* Sort of works

* Works until we get to concatenated tiles

* More attempts that don't work

* One bug down

* It made a tileset!

* Handle nonzero initial zooms

* Fix seeking within compressed feature streams

* Tests pass!

* Remove debug spew

* Oops: remember to delete the temporary files so they don't hang around

* Test that fails with the current compression code

* Properly account for bytes read while closing the compressed stream

* Limit the number of warnings about bad label points

* A little more armor when closing decompression

* This time for sure

* A different, less fragile, test that failed previously with compression

* Move feature stream compression to its own file

* Remove now-unused code to deserialize from a file

* Forgot to add the new files

* Remove a little debugging logging

* Add a couple of comments on what it means to be within decompression

* Fix indentation

* Update changelog. Remove stray debugging comment.
This commit is contained in:
Erica Fischer
2023-02-14 12:47:40 -08:00
committed by GitHub
parent b155b4671b
commit f2127fec97
23 changed files with 85678 additions and 455 deletions
+7
View File
@@ -1,3 +1,10 @@
## 2.23.0
* Remove the concept of "separate metadata." Features now always directly reference their keys and values rather than going through a second level of indirection.
* Limit the size of the string pool to 1/5 the size of memory, to prevent thrashing during feature ingestion.
* Avoid using writeable memory maps. Instead, explicitly copy data in and out of memory.
* Compress streams of features in the temporary files, to reduce disk usage and I/O latency
## 2.22.0 ## 2.22.0
* Speed up feature dropping by removing unnecessary search for other small features * Speed up feature dropping by removing unnecessary search for other small features
+2 -2
View File
@@ -47,7 +47,7 @@ C = $(wildcard *.c) $(wildcard *.cpp)
INCLUDES = -I/usr/local/include -I. INCLUDES = -I/usr/local/include -I.
LIBS = -L/usr/local/lib LIBS = -L/usr/local/lib
tippecanoe: geojson.o jsonpull/jsonpull.o tile.o pool.o mbtiles.o geometry.o projection.o memfile.o mvt.o serial.o main.o text.o dirtiles.o pmtiles_file.o plugin.o read_json.o write_json.o geobuf.o flatgeobuf.o evaluator.o geocsv.o csv.o geojson-loop.o json_logger.o visvalingam.o tippecanoe: geojson.o jsonpull/jsonpull.o tile.o pool.o mbtiles.o geometry.o projection.o memfile.o mvt.o serial.o main.o text.o dirtiles.o pmtiles_file.o plugin.o read_json.o write_json.o geobuf.o flatgeobuf.o evaluator.o geocsv.o csv.o geojson-loop.o json_logger.o visvalingam.o compression.o
$(CXX) $(PG) $(LIBS) $(FINAL_FLAGS) $(CXXFLAGS) -o $@ $^ $(LDFLAGS) -lm -lz -lsqlite3 -lpthread $(CXX) $(PG) $(LIBS) $(FINAL_FLAGS) $(CXXFLAGS) -o $@ $^ $(LDFLAGS) -lm -lz -lsqlite3 -lpthread
tippecanoe-enumerate: enumerate.o tippecanoe-enumerate: enumerate.o
@@ -56,7 +56,7 @@ tippecanoe-enumerate: enumerate.o
tippecanoe-decode: decode.o projection.o mvt.o write_json.o text.o jsonpull/jsonpull.o dirtiles.o pmtiles_file.o tippecanoe-decode: decode.o projection.o mvt.o write_json.o text.o jsonpull/jsonpull.o dirtiles.o pmtiles_file.o
$(CXX) $(PG) $(LIBS) $(FINAL_FLAGS) $(CXXFLAGS) -o $@ $^ $(LDFLAGS) -lm -lz -lsqlite3 $(CXX) $(PG) $(LIBS) $(FINAL_FLAGS) $(CXXFLAGS) -o $@ $^ $(LDFLAGS) -lm -lz -lsqlite3
tile-join: tile-join.o projection.o pool.o mbtiles.o mvt.o memfile.o dirtiles.o jsonpull/jsonpull.o text.o evaluator.o csv.o write_json.o pmtiles_file.o tile-join: tile-join.o projection.o mbtiles.o mvt.o memfile.o dirtiles.o jsonpull/jsonpull.o text.o evaluator.o csv.o write_json.o pmtiles_file.o
$(CXX) $(PG) $(LIBS) $(FINAL_FLAGS) $(CXXFLAGS) -o $@ $^ $(LDFLAGS) -lm -lz -lsqlite3 -lpthread $(CXX) $(PG) $(LIBS) $(FINAL_FLAGS) $(CXXFLAGS) -o $@ $^ $(LDFLAGS) -lm -lz -lsqlite3 -lpthread
tippecanoe-json-tool: jsontool.o jsonpull/jsonpull.o csv.o text.o geojson-loop.o tippecanoe-json-tool: jsontool.o jsonpull/jsonpull.o csv.o text.o geojson-loop.o
+269
View File
@@ -0,0 +1,269 @@
#ifdef __APPLE__
#define _DARWIN_UNLIMITED_STREAMS
#endif
#include "compression.hpp"
#include "errors.hpp"
#include "protozero/varint.hpp"
#include "serial.hpp"
void decompressor::begin() {
within = true;
zs.zalloc = NULL;
zs.zfree = NULL;
zs.opaque = NULL;
zs.msg = (char *) "";
int d = inflateInit(&zs);
if (d != Z_OK) {
fprintf(stderr, "initialize decompression: %d %s\n", d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
}
int decompressor::fread(void *p, size_t size, size_t nmemb, std::atomic<long long> *geompos) {
zs.next_out = (Bytef *) p;
zs.avail_out = size * nmemb;
while (zs.avail_out > 0) {
if (zs.avail_in == 0) {
size_t n = ::fread((Bytef *) buf.c_str(), sizeof(char), buf.size(), fp);
if (n == 0) {
if (within) {
fprintf(stderr, "Reached EOF while decompressing\n");
exit(EXIT_IMPOSSIBLE);
} else {
break;
}
}
zs.next_in = (Bytef *) buf.c_str();
zs.avail_in = n;
}
size_t avail_before = zs.avail_in;
if (within) {
int d = inflate(&zs, Z_NO_FLUSH);
*geompos += avail_before - zs.avail_in;
if (d == Z_OK) {
// it made some progress
} else if (d == Z_STREAM_END) {
// it may have made some progress and now we are done
within = false;
break;
} else {
fprintf(stderr, "decompression error %d %s\n", d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
} else {
size_t n = std::min(zs.avail_in, zs.avail_out);
memcpy(zs.next_out, zs.next_in, n);
*geompos += n;
zs.avail_out -= n;
zs.avail_in -= n;
zs.next_out += n;
zs.next_in += n;
}
}
return (size * nmemb - zs.avail_out) / size;
}
void decompressor::end(std::atomic<long long> *geompos) {
// "within" means that we haven't received end-of-stream yet,
// so consume more compressed data until we get there.
// This can be necessary if the caller knows that it is at
// the end of the feature stream (because it got a 0-length
// feature) but the decompressor doesn't know yet.
if (within) {
while (true) {
if (zs.avail_in == 0) {
size_t n = ::fread((Bytef *) buf.c_str(), sizeof(char), buf.size(), fp);
zs.next_in = (Bytef *) buf.c_str();
zs.avail_in = n;
}
zs.avail_out = 0;
size_t avail_before = zs.avail_in;
int d = inflate(&zs, Z_NO_FLUSH);
*geompos += avail_before - zs.avail_in;
if (d == Z_STREAM_END) {
break;
}
if (d == Z_OK) {
continue;
}
fprintf(stderr, "decompression: got %d, not Z_STREAM_END\n", d);
exit(EXIT_IMPOSSIBLE);
}
within = false;
}
int d = inflateEnd(&zs);
if (d != Z_OK) {
fprintf(stderr, "end decompression: %d %s\n", d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
}
int decompressor::deserialize_ulong_long(unsigned long long *zigzag, std::atomic<long long> *geompos) {
*zigzag = 0;
int shift = 0;
while (1) {
char c;
if (fread(&c, sizeof(char), 1, geompos) != 1) {
return 0;
}
if ((c & 0x80) == 0) {
*zigzag |= ((unsigned long long) c) << shift;
shift += 7;
break;
} else {
*zigzag |= ((unsigned long long) (c & 0x7F)) << shift;
shift += 7;
}
}
return 1;
}
int decompressor::deserialize_long_long(long long *n, std::atomic<long long> *geompos) {
unsigned long long zigzag = 0;
int ret = deserialize_ulong_long(&zigzag, geompos);
*n = protozero::decode_zigzag64(zigzag);
return ret;
}
int decompressor::deserialize_int(int *n, std::atomic<long long> *geompos) {
long long ll = 0;
int ret = deserialize_long_long(&ll, geompos);
*n = ll;
return ret;
}
int decompressor::deserialize_uint(unsigned *n, std::atomic<long long> *geompos) {
unsigned long long v;
deserialize_ulong_long(&v, geompos);
*n = v;
return 1;
}
void compressor::begin() {
zs.zalloc = NULL;
zs.zfree = NULL;
zs.opaque = NULL;
zs.msg = (char *) "";
int d = deflateInit(&zs, Z_DEFAULT_COMPRESSION);
if (d != Z_OK) {
fprintf(stderr, "initialize compression: %d %s\n", d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
}
void compressor::compressor::end(std::atomic<long long> *fpos, const char *fname) {
std::string buf;
buf.resize(5000);
if (zs.avail_in != 0) {
fprintf(stderr, "compression end called with data available\n");
exit(EXIT_IMPOSSIBLE);
}
zs.next_in = (Bytef *) buf.c_str();
zs.avail_in = 0;
while (true) {
zs.next_out = (Bytef *) buf.c_str();
zs.avail_out = buf.size();
int d = deflate(&zs, Z_FINISH);
::fwrite_check(buf.c_str(), sizeof(char), zs.next_out - (Bytef *) buf.c_str(), fp, fpos, fname);
if (d == Z_OK || d == Z_BUF_ERROR) {
// it can take several calls to flush out all the buffered data
continue;
}
if (d != Z_STREAM_END) {
fprintf(stderr, "%s: finish compression: %d %s\n", fname, d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
break;
}
zs.next_out = (Bytef *) buf.c_str();
zs.avail_out = buf.size();
int d = deflateEnd(&zs);
if (d != Z_OK) {
fprintf(stderr, "%s: end compression: %d %s\n", fname, d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
::fwrite_check(buf.c_str(), sizeof(char), zs.next_out - (Bytef *) buf.c_str(), fp, fpos, fname);
}
int compressor::fclose() {
return ::fclose(fp);
}
void compressor::fwrite_check(const char *p, size_t size, size_t nmemb, std::atomic<long long> *fpos, const char *fname) {
std::string buf;
buf.resize(size * nmemb * 2 + 200);
zs.next_in = (Bytef *) p;
zs.avail_in = size * nmemb;
while (zs.avail_in > 0) {
zs.next_out = (Bytef *) buf.c_str();
zs.avail_out = buf.size();
int d = deflate(&zs, Z_NO_FLUSH);
if (d != Z_OK) {
fprintf(stderr, "%s: deflate: %d %s\n", fname, d, zs.msg);
exit(EXIT_IMPOSSIBLE);
}
::fwrite_check(buf.c_str(), sizeof(char), zs.next_out - (Bytef *) buf.c_str(), fp, fpos, fname);
}
}
void compressor::serialize_ulong_long(unsigned long long val, std::atomic<long long> *fpos, const char *fname) {
while (1) {
unsigned char b = val & 0x7F;
if ((val >> 7) != 0) {
b |= 0x80;
fwrite_check((const char *) &b, 1, 1, fpos, fname);
val >>= 7;
} else {
fwrite_check((const char *) &b, 1, 1, fpos, fname);
break;
}
}
}
void compressor::serialize_long_long(long long val, std::atomic<long long> *fpos, const char *fname) {
unsigned long long zigzag = protozero::encode_zigzag64(val);
serialize_ulong_long(zigzag, fpos, fname);
}
void compressor::serialize_int(int val, std::atomic<long long> *fpos, const char *fname) {
serialize_long_long(val, fpos, fname);
}
void compressor::serialize_uint(unsigned val, std::atomic<long long> *fpos, const char *fname) {
serialize_ulong_long(val, fpos, fname);
}
+58
View File
@@ -0,0 +1,58 @@
#ifdef __APPLE__
#define _DARWIN_UNLIMITED_STREAMS
#endif
#include <stdio.h>
#include <string>
#include <atomic>
#include <zlib.h>
struct decompressor {
FILE *fp = NULL;
z_stream zs;
std::string buf;
// from begin() to receiving end-of-stream
bool within = false;
decompressor(FILE *f) {
fp = f;
buf.resize(5000);
zs.next_in = (Bytef *) buf.c_str();
zs.avail_in = 0;
}
decompressor() {
}
void begin();
int fread(void *p, size_t size, size_t nmemb, std::atomic<long long> *geompos);
void end(std::atomic<long long> *geompos);
int deserialize_ulong_long(unsigned long long *zigzag, std::atomic<long long> *geompos);
int deserialize_long_long(long long *n, std::atomic<long long> *geompos);
int deserialize_int(int *n, std::atomic<long long> *geompos);
int deserialize_uint(unsigned *n, std::atomic<long long> *geompos);
};
struct compressor {
FILE *fp = NULL;
z_stream zs;
compressor(FILE *f) {
fp = f;
}
compressor() {
}
void begin();
void end(std::atomic<long long> *fpos, const char *fname);
int fclose();
void fwrite_check(const char *p, size_t size, size_t nmemb, std::atomic<long long> *fpos, const char *fname);
void serialize_ulong_long(unsigned long long val, std::atomic<long long> *fpos, const char *fname);
void serialize_long_long(long long val, std::atomic<long long> *fpos, const char *fname);
void serialize_int(int val, std::atomic<long long> *fpos, const char *fname);
void serialize_uint(unsigned val, std::atomic<long long> *fpos, const char *fname);
};
+11 -9
View File
@@ -25,7 +25,7 @@
static int clip(double *x0, double *y0, double *x1, double *y1, double xmin, double ymin, double xmax, double ymax); static int clip(double *x0, double *y0, double *x1, double *y1, double xmin, double ymin, double xmax, double ymax);
drawvec decode_geometry(FILE *meta, std::atomic<long long> *geompos, int z, unsigned tx, unsigned ty, long long *bbox, unsigned initial_x, unsigned initial_y) { drawvec decode_geometry(char **meta, int z, unsigned tx, unsigned ty, long long *bbox, unsigned initial_x, unsigned initial_y) {
drawvec out; drawvec out;
bbox[0] = LLONG_MAX; bbox[0] = LLONG_MAX;
@@ -38,10 +38,7 @@ drawvec decode_geometry(FILE *meta, std::atomic<long long> *geompos, int z, unsi
while (1) { while (1) {
draw d; draw d;
if (!deserialize_byte_io(meta, &d.op, geompos)) { deserialize_byte(meta, &d.op);
fprintf(stderr, "Internal error: Unexpected end of file in geometry\n");
exit(EXIT_IMPOSSIBLE);
}
if (d.op == VT_END) { if (d.op == VT_END) {
break; break;
} }
@@ -49,8 +46,8 @@ drawvec decode_geometry(FILE *meta, std::atomic<long long> *geompos, int z, unsi
if (d.op == VT_MOVETO || d.op == VT_LINETO) { if (d.op == VT_MOVETO || d.op == VT_LINETO) {
long long dx, dy; long long dx, dy;
deserialize_long_long_io(meta, &dx, geompos); deserialize_long_long(meta, &dx);
deserialize_long_long_io(meta, &dy, geompos); deserialize_long_long(meta, &dy);
wx += dx * (1 << geometry_scale); wx += dx * (1 << geometry_scale);
wy += dy * (1 << geometry_scale); wy += dy * (1 << geometry_scale);
@@ -1541,7 +1538,8 @@ struct sorty {
struct sorty_sorter { struct sorty_sorter {
int kind; int kind;
sorty_sorter(int k) : kind(k) {}; sorty_sorter(int k)
: kind(k){};
bool operator()(const sorty &a, const sorty &b) const { bool operator()(const sorty &a, const sorty &b) const {
long long xa, ya, xb, yb; long long xa, ya, xb, yb;
@@ -1775,7 +1773,11 @@ drawvec polygon_to_anchor(const drawvec &geom) {
if (goodness <= 0) { if (goodness <= 0) {
double lon, lat; double lon, lat;
tile2lonlat(d.x, d.y, 32, &lon, &lat); tile2lonlat(d.x, d.y, 32, &lon, &lat);
fprintf(stderr, "could not find label point: %s %f,%f\n", kind, lat, lon);
static std::atomic<long long> warned(0);
if (warned++ < 10) {
fprintf(stderr, "could not find good label point: %s %f,%f\n", kind, lat, lon);
}
} }
} }
+1 -1
View File
@@ -58,7 +58,7 @@ struct draw {
typedef std::vector<draw> drawvec; typedef std::vector<draw> drawvec;
struct serial_feature; struct serial_feature;
drawvec decode_geometry(FILE *meta, std::atomic<long long> *geompos, int z, unsigned tx, unsigned ty, long long *bbox, unsigned initial_x, unsigned initial_y); drawvec decode_geometry(char **meta, int z, unsigned tx, unsigned ty, long long *bbox, unsigned initial_x, unsigned initial_y);
void to_tile_scale(drawvec &geom, int z, int detail); void to_tile_scale(drawvec &geom, int z, int detail);
drawvec from_tile_scale(drawvec const &geom, int z, int detail); drawvec from_tile_scale(drawvec const &geom, int z, int detail);
drawvec remove_noop(drawvec geom, int type, int shift); drawvec remove_noop(drawvec geom, int type, int shift);
+98 -139
View File
@@ -105,6 +105,7 @@ struct source {
size_t CPUS; size_t CPUS;
size_t TEMP_FILES; size_t TEMP_FILES;
long long MAX_FILES; long long MAX_FILES;
size_t memsize;
static long long diskfree; static long long diskfree;
char **av; char **av;
@@ -113,9 +114,9 @@ std::vector<clipbbox> clipbboxes;
void checkdisk(std::vector<struct reader> *r) { void checkdisk(std::vector<struct reader> *r) {
long long used = 0; long long used = 0;
for (size_t i = 0; i < r->size(); i++) { for (size_t i = 0; i < r->size(); i++) {
// Meta, pool, and tree are used once. // Pool and tree are used once.
// Geometry and index will be duplicated during sorting and tiling. // Geometry and index will be duplicated during sorting and tiling.
used += (*r)[i].metapos + 2 * (*r)[i].geompos + 2 * (*r)[i].indexpos + (*r)[i].poolfile->len + (*r)[i].treefile->len; used += 2 * (*r)[i].geompos + 2 * (*r)[i].indexpos + (*r)[i].poolfile->off + (*r)[i].treefile->off;
} }
static int warned = 0; static int warned = 0;
@@ -326,8 +327,11 @@ static void merge(struct mergelist *merges, size_t nmerges, unsigned char *map,
while (head != NULL) { while (head != NULL) {
struct index ix = *((struct index *) (map + head->start)); struct index ix = *((struct index *) (map + head->start));
long long pos = *geompos; long long pos = *geompos;
fwrite_check(geom_map + ix.start, 1, ix.end - ix.start, geom_out, "merge geometry");
*geompos += ix.end - ix.start; // MAGIC: This knows that the feature minzoom is the last byte of the serialized feature
// and is writing one byte less and then adding the byte for the minzoom.
fwrite_check(geom_map + ix.start, 1, ix.end - ix.start - 1, geom_out, geompos, "merge geometry");
int feature_minzoom = calc_feature_minzoom(&ix, ds, maxzoom, gamma); int feature_minzoom = calc_feature_minzoom(&ix, ds, maxzoom, gamma);
serialize_byte(geom_out, feature_minzoom, geompos, "merge geometry"); serialize_byte(geom_out, feature_minzoom, geompos, "merge geometry");
@@ -335,12 +339,14 @@ static void merge(struct mergelist *merges, size_t nmerges, unsigned char *map,
*progress += (ix.end - ix.start) * 3 / 4; *progress += (ix.end - ix.start) * 3 / 4;
if (!quiet && !quiet_progress && progress_time() && 100 * *progress / *progress_max != *progress_reported) { if (!quiet && !quiet_progress && progress_time() && 100 * *progress / *progress_max != *progress_reported) {
fprintf(stderr, "Reordering geometry: %lld%% \r", 100 * *progress / *progress_max); fprintf(stderr, "Reordering geometry: %lld%% \r", 100 * *progress / *progress_max);
fflush(stderr);
*progress_reported = 100 * *progress / *progress_max; *progress_reported = 100 * *progress / *progress_max;
} }
ix.start = pos; ix.start = pos;
ix.end = *geompos; ix.end = *geompos;
fwrite_check(&ix, bytes, 1, indexfile, "merge temporary"); std::atomic<long long> indexpos;
fwrite_check(&ix, bytes, 1, indexfile, &indexpos, "merge temporary");
head->start += bytes; head->start += bytes;
struct mergelist *m = head; struct mergelist *m = head;
@@ -382,32 +388,25 @@ void *run_sort(void *v) {
a->merges[start / a->unit].end = end; a->merges[start / a->unit].end = end;
a->merges[start / a->unit].next = NULL; a->merges[start / a->unit].next = NULL;
// MAP_PRIVATE to avoid disk writes if it fits in memory // Read section of index into memory to sort and then use pwrite()
void *map = mmap(NULL, end - start, PROT_READ | PROT_WRITE, MAP_PRIVATE, a->indexfd, start); // to write it back out rather than sorting in mapped memory,
if (map == MAP_FAILED) { // because writable mapped memory seems to have bad performance
perror("mmap in run_sort"); // problems on ECS (and maybe in containers in general)?
exit(EXIT_MEMORY);
std::string s;
s.resize(end - start);
if (pread(a->indexfd, (void *) s.c_str(), end - start, start) != end - start) {
fprintf(stderr, "pread(index): %s\n", strerror(errno));
exit(EXIT_READ);
} }
madvise(map, end - start, MADV_RANDOM);
madvise(map, end - start, MADV_WILLNEED);
qsort(map, (end - start) / a->bytes, a->bytes, indexcmp); qsort((void *) s.c_str(), (end - start) / a->bytes, a->bytes, indexcmp);
// Sorting and then copying avoids disk access to if (pwrite(a->indexfd, s.c_str(), end - start, start) != end - start) {
// write out intermediate stages of the sort. fprintf(stderr, "pwrite(index): %s\n", strerror(errno));
exit(EXIT_WRITE);
void *map2 = mmap(NULL, end - start, PROT_READ | PROT_WRITE, MAP_SHARED, a->indexfd, start);
if (map2 == MAP_FAILED) {
perror("mmap (write)");
exit(EXIT_MEMORY);
} }
madvise(map2, end - start, MADV_SEQUENTIAL);
memcpy(map2, map, end - start);
// No madvise, since caller will want the sorted data
munmap(map, end - start);
munmap(map2, end - start);
} }
return NULL; return NULL;
@@ -788,20 +787,21 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split
unsigned long long which = (ix.ix << prefix) >> (64 - splitbits); unsigned long long which = (ix.ix << prefix) >> (64 - splitbits);
long long pos = sub_geompos[which]; long long pos = sub_geompos[which];
fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfiles[which], "geom"); fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfiles[which], &sub_geompos[which], "geom");
sub_geompos[which] += ix.end - ix.start;
// Count this as a 25%-accomplishment, since we will copy again // Count this as a 25%-accomplishment, since we will copy again
*progress += (ix.end - ix.start) / 4; *progress += (ix.end - ix.start) / 4;
if (!quiet && !quiet_progress && progress_time() && 100 * *progress / *progress_max != *progress_reported) { if (!quiet && !quiet_progress && progress_time() && 100 * *progress / *progress_max != *progress_reported) {
fprintf(stderr, "Reordering geometry: %lld%% \r", 100 * *progress / *progress_max); fprintf(stderr, "Reordering geometry: %lld%% \r", 100 * *progress / *progress_max);
fflush(stderr);
*progress_reported = 100 * *progress / *progress_max; *progress_reported = 100 * *progress / *progress_max;
} }
ix.start = pos; ix.start = pos;
ix.end = sub_geompos[which]; ix.end = sub_geompos[which];
fwrite_check(&ix, sizeof(struct index), 1, indexfiles[which], "index"); std::atomic<long long> indexpos;
fwrite_check(&ix, sizeof(struct index), 1, indexfiles[which], &indexpos, "index");
} }
madvise(indexmap, indexst.st_size, MADV_DONTNEED); madvise(indexmap, indexst.st_size, MADV_DONTNEED);
@@ -958,8 +958,7 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split
struct index ix = indexmap[a]; struct index ix = indexmap[a];
long long pos = *geompos_out; long long pos = *geompos_out;
fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfile, "geom"); fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfile, geompos_out, "geom");
*geompos_out += ix.end - ix.start;
int feature_minzoom = calc_feature_minzoom(&ix, ds, maxzoom, gamma); int feature_minzoom = calc_feature_minzoom(&ix, ds, maxzoom, gamma);
serialize_byte(geomfile, feature_minzoom, geompos_out, "merge geometry"); serialize_byte(geomfile, feature_minzoom, geompos_out, "merge geometry");
@@ -967,12 +966,14 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split
*progress += (ix.end - ix.start) * 3 / 4; *progress += (ix.end - ix.start) * 3 / 4;
if (!quiet && !quiet_progress && progress_time() && 100 * *progress / *progress_max != *progress_reported) { if (!quiet && !quiet_progress && progress_time() && 100 * *progress / *progress_max != *progress_reported) {
fprintf(stderr, "Reordering geometry: %lld%% \r", 100 * *progress / *progress_max); fprintf(stderr, "Reordering geometry: %lld%% \r", 100 * *progress / *progress_max);
fflush(stderr);
*progress_reported = 100 * *progress / *progress_max; *progress_reported = 100 * *progress / *progress_max;
} }
ix.start = pos; ix.start = pos;
ix.end = *geompos_out; ix.end = *geompos_out;
fwrite_check(&ix, sizeof(struct index), 1, indexfile, "index"); std::atomic<long long> indexpos;
fwrite_check(&ix, sizeof(struct index), 1, indexfile, &indexpos, "index");
} }
madvise(indexmap, indexst.st_size, MADV_DONTNEED); madvise(indexmap, indexst.st_size, MADV_DONTNEED);
@@ -1028,17 +1029,8 @@ void prep_drop_states(struct drop_state *ds, int maxzoom, int basezoom, double d
} }
} }
void radix(std::vector<struct reader> &readers, int nreaders, FILE *geomfile, FILE *indexfile, const char *tmpdir, std::atomic<long long> *geompos, int maxzoom, int basezoom, double droprate, double gamma) { static size_t calc_memsize() {
// Run through the index and geometry for each reader, size_t mem;
// splitting the contents out by index into as many
// sub-files as we can write to simultaneously.
// Then sort each of those by index, recursively if it is
// too big to fit in memory.
// Then concatenate each of the sub-outputs into a final output.
long long mem;
#ifdef __APPLE__ #ifdef __APPLE__
int64_t hw_memsize; int64_t hw_memsize;
@@ -1059,6 +1051,21 @@ void radix(std::vector<struct reader> &readers, int nreaders, FILE *geomfile, FI
mem = (long long) pages * pagesize; mem = (long long) pages * pagesize;
#endif #endif
return mem;
}
void radix(std::vector<struct reader> &readers, int nreaders, FILE *geomfile, FILE *indexfile, const char *tmpdir, std::atomic<long long> *geompos, int maxzoom, int basezoom, double droprate, double gamma) {
// Run through the index and geometry for each reader,
// splitting the contents out by index into as many
// sub-files as we can write to simultaneously.
// Then sort each of those by index, recursively if it is
// too big to fit in memory.
// Then concatenate each of the sub-outputs into a final output.
long long mem = memsize;
// Just for code coverage testing. Deeply recursive sorting is very slow // Just for code coverage testing. Deeply recursive sorting is very slow
// compared to sorting in memory. // compared to sorting in memory.
if (additional[A_PREFER_RADIX_SORT]) { if (additional[A_PREFER_RADIX_SORT]) {
@@ -1066,7 +1073,7 @@ void radix(std::vector<struct reader> &readers, int nreaders, FILE *geomfile, FI
} }
long long availfiles = MAX_FILES - 2 * nreaders // each reader has a geom and an index long long availfiles = MAX_FILES - 2 * nreaders // each reader has a geom and an index
- 4 // pool, meta, mbtiles, mbtiles journal - 3 // pool, mbtiles, mbtiles journal
- 4 // top-level geom and index output, both FILE and fd - 4 // top-level geom and index output, both FILE and fd
- 3; // stdin, stdout, stderr - 3; // stdin, stdout, stderr
@@ -1164,23 +1171,16 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
for (size_t i = 0; i < CPUS; i++) { for (size_t i = 0; i < CPUS; i++) {
struct reader *r = &readers[i]; struct reader *r = &readers[i];
char metaname[strlen(tmpdir) + strlen("/meta.XXXXXXXX") + 1];
char poolname[strlen(tmpdir) + strlen("/pool.XXXXXXXX") + 1]; char poolname[strlen(tmpdir) + strlen("/pool.XXXXXXXX") + 1];
char treename[strlen(tmpdir) + strlen("/tree.XXXXXXXX") + 1]; char treename[strlen(tmpdir) + strlen("/tree.XXXXXXXX") + 1];
char geomname[strlen(tmpdir) + strlen("/geom.XXXXXXXX") + 1]; char geomname[strlen(tmpdir) + strlen("/geom.XXXXXXXX") + 1];
char indexname[strlen(tmpdir) + strlen("/index.XXXXXXXX") + 1]; char indexname[strlen(tmpdir) + strlen("/index.XXXXXXXX") + 1];
sprintf(metaname, "%s%s", tmpdir, "/meta.XXXXXXXX");
sprintf(poolname, "%s%s", tmpdir, "/pool.XXXXXXXX"); sprintf(poolname, "%s%s", tmpdir, "/pool.XXXXXXXX");
sprintf(treename, "%s%s", tmpdir, "/tree.XXXXXXXX"); sprintf(treename, "%s%s", tmpdir, "/tree.XXXXXXXX");
sprintf(geomname, "%s%s", tmpdir, "/geom.XXXXXXXX"); sprintf(geomname, "%s%s", tmpdir, "/geom.XXXXXXXX");
sprintf(indexname, "%s%s", tmpdir, "/index.XXXXXXXX"); sprintf(indexname, "%s%s", tmpdir, "/index.XXXXXXXX");
r->metafd = mkstemp_cloexec(metaname);
if (r->metafd < 0) {
perror(metaname);
exit(EXIT_OPEN);
}
r->poolfd = mkstemp_cloexec(poolname); r->poolfd = mkstemp_cloexec(poolname);
if (r->poolfd < 0) { if (r->poolfd < 0) {
perror(poolname); perror(poolname);
@@ -1202,11 +1202,6 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
exit(EXIT_OPEN); exit(EXIT_OPEN);
} }
r->metafile = fopen_oflag(metaname, "wb", O_WRONLY | O_CLOEXEC);
if (r->metafile == NULL) {
perror(metaname);
exit(EXIT_OPEN);
}
r->poolfile = memfile_open(r->poolfd); r->poolfile = memfile_open(r->poolfd);
if (r->poolfile == NULL) { if (r->poolfile == NULL) {
perror(poolname); perror(poolname);
@@ -1227,11 +1222,9 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
perror(indexname); perror(indexname);
exit(EXIT_OPEN); exit(EXIT_OPEN);
} }
r->metapos = 0;
r->geompos = 0; r->geompos = 0;
r->indexpos = 0; r->indexpos = 0;
unlink(metaname);
unlink(poolname); unlink(poolname);
unlink(treename); unlink(treename);
unlink(geomname); unlink(geomname);
@@ -1242,8 +1235,6 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
struct stringpool p; struct stringpool p;
memfile_write(r->treefile, &p, sizeof(struct stringpool)); memfile_write(r->treefile, &p, sizeof(struct stringpool));
} }
// Keep metadata file from being completely empty if no attributes
serialize_int(r->metafile, 0, &r->metapos, "meta");
r->file_bbox[0] = r->file_bbox[1] = UINT_MAX; r->file_bbox[0] = r->file_bbox[1] = UINT_MAX;
r->file_bbox[2] = r->file_bbox[3] = 0; r->file_bbox[2] = r->file_bbox[3] = 0;
@@ -1674,7 +1665,8 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
int n; int n;
while ((n = fp->read(buf, READ_BUF)) > 0) { while ((n = fp->read(buf, READ_BUF)) > 0) {
fwrite_check(buf, sizeof(char), n, readfp, reading.c_str()); std::atomic<long long> readingpos;
fwrite_check(buf, sizeof(char), n, readfp, &readingpos, reading.c_str());
ahead += n; ahead += n;
if (buf[n - 1] == read_parallel_this && ahead > PARSE_MIN) { if (buf[n - 1] == read_parallel_this && ahead > PARSE_MIN) {
@@ -1803,13 +1795,10 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
if (!quiet) { if (!quiet) {
fprintf(stderr, " \r"); fprintf(stderr, " \r");
// (stderr, "Read 10000.00 million features\r", *progress_seq / 1000000.0); // (stderr, "Read 10000.00 million features\r", *progress_seq / 1000000.0);
fflush(stderr);
} }
for (size_t i = 0; i < CPUS; i++) { for (size_t i = 0; i < CPUS; i++) {
if (fclose(readers[i].metafile) != 0) {
perror("fclose meta");
exit(EXIT_CLOSE);
}
if (fclose(readers[i].geomfile) != 0) { if (fclose(readers[i].geomfile) != 0) {
perror("fclose geom"); perror("fclose geom");
exit(EXIT_CLOSE); exit(EXIT_CLOSE);
@@ -1824,21 +1813,16 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
perror("stat geom\n"); perror("stat geom\n");
exit(EXIT_STAT); exit(EXIT_STAT);
} }
if (fstat(readers[i].metafd, &readers[i].metast) != 0) {
perror("stat meta\n");
exit(EXIT_STAT);
}
} }
// Create a combined string pool and a combined metadata file // Create a combined string pool
// but keep track of the offsets into it since we still need // but keep track of the offsets into it since we still need
// segment+offset to find the data. // segment+offset to find the data.
// 2 * CPUS: One per input thread, one per tiling thread // 2 * CPUS: One per input thread, one per tiling thread
long long pool_off[2 * CPUS]; long long pool_off[2 * CPUS];
long long meta_off[2 * CPUS];
for (size_t i = 0; i < 2 * CPUS; i++) { for (size_t i = 0; i < 2 * CPUS; i++) {
pool_off[i] = meta_off[i] = 0; pool_off[i] = 0;
} }
char poolname[strlen(tmpdir) + strlen("/pool.XXXXXXXX") + 1]; char poolname[strlen(tmpdir) + strlen("/pool.XXXXXXXX") + 1];
@@ -1858,60 +1842,51 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
unlink(poolname); unlink(poolname);
char metaname[strlen(tmpdir) + strlen("/meta.XXXXXXXX") + 1];
sprintf(metaname, "%s%s", tmpdir, "/meta.XXXXXXXX");
int metafd = mkstemp_cloexec(metaname);
if (metafd < 0) {
perror(metaname);
exit(EXIT_OPEN);
}
FILE *metafile = fopen_oflag(metaname, "wb", O_WRONLY | O_CLOEXEC);
if (metafile == NULL) {
perror(metaname);
exit(EXIT_OPEN);
}
unlink(metaname);
std::atomic<long long> metapos(0);
std::atomic<long long> poolpos(0); std::atomic<long long> poolpos(0);
for (size_t i = 0; i < CPUS; i++) { for (size_t i = 0; i < CPUS; i++) {
if (readers[i].metapos > 0) { // If the memfile is not done yet, it is in memory, so just copy the memory.
void *map = mmap(NULL, readers[i].metapos, PROT_READ, MAP_PRIVATE, readers[i].metafd, 0); // Otherwise, we need to merge memory and file.
if (map == MAP_FAILED) {
perror("mmap unmerged meta");
exit(EXIT_MEMORY);
}
madvise(map, readers[i].metapos, MADV_SEQUENTIAL);
madvise(map, readers[i].metapos, MADV_WILLNEED);
if (fwrite(map, readers[i].metapos, 1, metafile) != 1) {
perror("Reunify meta");
exit(EXIT_WRITE);
}
madvise(map, readers[i].metapos, MADV_DONTNEED);
if (munmap(map, readers[i].metapos) != 0) {
perror("unmap unmerged meta");
}
}
meta_off[i] = metapos; if (readers[i].poolfile->fp == NULL) {
metapos += readers[i].metapos; // still in memory
if (close(readers[i].metafd) != 0) {
perror("close unmerged meta");
}
if (readers[i].poolfile->off > 0) { if (readers[i].poolfile->map.size() > 0) {
if (fwrite(readers[i].poolfile->map, readers[i].poolfile->off, 1, poolfile) != 1) { if (fwrite(readers[i].poolfile->map.c_str(), readers[i].poolfile->map.size(), 1, poolfile) != 1) {
perror("Reunify string pool"); perror("Reunify string pool");
exit(EXIT_WRITE); exit(EXIT_WRITE);
} }
} }
pool_off[i] = poolpos;
poolpos += readers[i].poolfile->map.size();
} else {
// split into memory and file
if (fflush(readers[i].poolfile->fp) != 0) {
perror("fflush poolfile");
exit(EXIT_WRITE);
}
char *s = (char *) mmap(NULL, readers[i].poolfile->off, PROT_READ, MAP_PRIVATE, readers[i].poolfile->fd, 0);
if (s == MAP_FAILED) {
perror("mmap string pool for copy");
exit(EXIT_MEMORY);
}
madvise(s, readers[i].poolfile->off, MADV_SEQUENTIAL);
if (fwrite(s, sizeof(char), readers[i].poolfile->off, poolfile) != readers[i].poolfile->off) {
perror("Reunify string pool (split)");
exit(EXIT_WRITE);
}
if (munmap(s, readers[i].poolfile->off) != 0) {
perror("unmap string pool for copy");
exit(EXIT_MEMORY);
}
pool_off[i] = poolpos; pool_off[i] = poolpos;
poolpos += readers[i].poolfile->off; poolpos += readers[i].poolfile->off;
}
memfile_close(readers[i].poolfile); memfile_close(readers[i].poolfile);
} }
@@ -1919,17 +1894,6 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
perror("fclose pool"); perror("fclose pool");
exit(EXIT_CLOSE); exit(EXIT_CLOSE);
} }
if (fclose(metafile) != 0) {
perror("fclose meta");
exit(EXIT_CLOSE);
}
char *meta = (char *) mmap(NULL, metapos, PROT_READ, MAP_PRIVATE, metafd, 0);
if (meta == MAP_FAILED) {
perror("mmap meta");
exit(EXIT_MEMORY);
}
madvise(meta, metapos, MADV_RANDOM);
char *stringpool = NULL; char *stringpool = NULL;
if (poolpos > 0) { // Will be 0 if -X was specified if (poolpos > 0) { // Will be 0 if -X was specified
@@ -1983,7 +1947,7 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
std::atomic<long long> geompos(0); std::atomic<long long> geompos(0);
/* initial tile is 0/0/0 */ /* initial tile is normally 0/0/0 but can be iz/ix/iy if limited to one tile */
serialize_int(geomfile, iz, &geompos, fname); serialize_int(geomfile, iz, &geompos, fname);
serialize_uint(geomfile, ix, &geompos, fname); serialize_uint(geomfile, ix, &geompos, fname);
serialize_uint(geomfile, iy, &geompos, fname); serialize_uint(geomfile, iy, &geompos, fname);
@@ -1991,7 +1955,7 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
radix(readers, CPUS, geomfile, indexfile, tmpdir, &geompos, maxzoom, basezoom, droprate, gamma); radix(readers, CPUS, geomfile, indexfile, tmpdir, &geompos, maxzoom, basezoom, droprate, gamma);
/* end of tile */ /* end of tile */
serialize_byte(geomfile, -2, &geompos, fname); serialize_ulong_long(geomfile, 0, &geompos, fname); // EOF
if (fclose(geomfile) != 0) { if (fclose(geomfile) != 0) {
perror("fclose geom"); perror("fclose geom");
@@ -2014,9 +1978,8 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
if (!quiet) { if (!quiet) {
long long s = progress_seq; long long s = progress_seq;
long long geompos_print = geompos; long long geompos_print = geompos;
long long metapos_print = metapos;
long long poolpos_print = poolpos; long long poolpos_print = poolpos;
fprintf(stderr, "%lld features, %lld bytes of geometry, %lld bytes of separate metadata, %lld bytes of string pool\n", s, geompos_print, metapos_print, poolpos_print); fprintf(stderr, "%lld features, %lld bytes of geometry, %lld bytes of string pool\n", s, geompos_print, poolpos_print);
} }
if (indexpos == 0) { if (indexpos == 0) {
@@ -2083,6 +2046,7 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
progress = nprogress; progress = nprogress;
if (!quiet && !quiet_progress && progress_time()) { if (!quiet && !quiet_progress && progress_time()) {
fprintf(stderr, "Maxzoom: %lld%% \r", progress); fprintf(stderr, "Maxzoom: %lld%% \r", progress);
fflush(stderr);
} }
} }
} }
@@ -2269,6 +2233,7 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
progress = nprogress; progress = nprogress;
if (!quiet && !quiet_progress && progress_time()) { if (!quiet && !quiet_progress && progress_time()) {
fprintf(stderr, "Base zoom/drop rate: %lld%% \r", progress); fprintf(stderr, "Base zoom/drop rate: %lld%% \r", progress);
fflush(stderr);
} }
} }
@@ -2518,7 +2483,7 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
std::atomic<unsigned> midx(0); std::atomic<unsigned> midx(0);
std::atomic<unsigned> midy(0); std::atomic<unsigned> midy(0);
std::vector<strategy> strategies; std::vector<strategy> strategies;
int written = traverse_zooms(fd, size, meta, stringpool, &midx, &midy, maxzoom, minzoom, outdb, outdir, buffer, fname, tmpdir, gamma, full_detail, low_detail, min_detail, meta_off, pool_off, initial_x, initial_y, simplification, maxzoom_simplification, layermaps, prefilter, postfilter, attribute_accum, filter, strategies); int written = traverse_zooms(fd, size, stringpool, &midx, &midy, maxzoom, minzoom, outdb, outdir, buffer, fname, tmpdir, gamma, full_detail, low_detail, min_detail, pool_off, initial_x, initial_y, simplification, maxzoom_simplification, layermaps, prefilter, postfilter, attribute_accum, filter, strategies, iz);
if (maxzoom != written) { if (maxzoom != written) {
if (written > minzoom) { if (written > minzoom) {
@@ -2531,14 +2496,6 @@ std::pair<int, metadata> read_input(std::vector<source> &sources, char *fname, i
} }
} }
madvise(meta, metapos, MADV_DONTNEED);
if (munmap(meta, metapos) != 0) {
perror("munmap meta");
}
if (close(metafd) < 0) {
perror("close meta");
}
if (poolpos > 0) { if (poolpos > 0) {
madvise((void *) stringpool, poolpos, MADV_DONTNEED); madvise((void *) stringpool, poolpos, MADV_DONTNEED);
if (munmap(stringpool, poolpos) != 0) { if (munmap(stringpool, poolpos) != 0) {
@@ -2745,6 +2702,8 @@ int main(int argc, char **argv) {
int files_open_at_start; int files_open_at_start;
json_object *filter = NULL; json_object *filter = NULL;
memsize = calc_memsize();
for (i = 0; i < 256; i++) { for (i = 0; i < 256; i++) {
prevent[i] = 0; prevent[i] = 0;
additional[i] = 0; additional[i] = 0;
+2
View File
@@ -4,6 +4,7 @@
#include <stddef.h> #include <stddef.h>
#include <atomic> #include <atomic>
#include <string> #include <string>
#include <vector>
#include "json_logger.hpp" #include "json_logger.hpp"
@@ -47,6 +48,7 @@ extern int extra_detail;
extern size_t CPUS; extern size_t CPUS;
extern size_t TEMP_FILES; extern size_t TEMP_FILES;
extern size_t memsize;
extern size_t max_tile_size; extern size_t max_tile_size;
extern size_t max_tile_features; extern size_t max_tile_features;
+43 -26
View File
@@ -12,28 +12,24 @@ struct memfile *memfile_open(int fd) {
return NULL; return NULL;
} }
char *map = (char *) mmap(NULL, INITIAL, PROT_READ | PROT_WRITE, MAP_SHARED, fd, 0);
if (map == MAP_FAILED) {
return NULL;
}
struct memfile *mf = new memfile; struct memfile *mf = new memfile;
if (mf == NULL) { if (mf == NULL) {
munmap(map, INITIAL);
return NULL; return NULL;
} }
mf->fd = fd; mf->fd = fd;
mf->map = map;
mf->len = INITIAL;
mf->off = 0;
mf->tree = 0; mf->tree = 0;
mf->off = 0;
return mf; return mf;
} }
int memfile_close(struct memfile *file) { int memfile_close(struct memfile *file) {
if (munmap(file->map, file->len) != 0) { // If it isn't full yet, flush out the string to the file now.
// If it is full, close out the buffered file writer.
if (file->fp == NULL) {
if (write(file->fd, file->map.c_str(), file->map.size()) != (ssize_t) file->map.size()) {
return -1; return -1;
} }
@@ -42,30 +38,51 @@ int memfile_close(struct memfile *file) {
return -1; return -1;
} }
} }
} else {
if (fclose(file->fp) != 0) {
return -1;
}
}
delete file; delete file;
return 0; return 0;
} }
int memfile_write(struct memfile *file, void *s, long long len) { int memfile_write(struct memfile *file, void *s, long long len) {
if (file->off + len > file->len) { // If it is full, append to the file.
if (munmap(file->map, file->len) != 0) { // If it is not full yet, append to the string in memory.
return -1;
}
file->len += (len + INCREMENT + 1) / INCREMENT * INCREMENT; if (file->fp != NULL) {
if (fwrite(s, sizeof(char), len, file->fp) != (size_t) len) {
if (ftruncate(file->fd, file->len) != 0) { return 0;
return -1;
} }
file->map = (char *) mmap(NULL, file->len, PROT_READ | PROT_WRITE, MAP_SHARED, file->fd, 0);
if (file->map == MAP_FAILED) {
return -1;
}
}
memcpy(file->map + file->off, s, len);
file->off += len; file->off += len;
} else {
file->map.append(std::string((char *) s, len));
file->off += len;
}
return len; return len;
} }
void memfile_full(struct memfile *file) {
// The file is full. Write out a copy of whatever has accumulated in memory
// to the file, and switch to appending to the file. Existing references
// into the memory still work.
if (file->fp != NULL) {
fprintf(stderr, "memfile marked full twice\n");
exit(EXIT_IMPOSSIBLE);
}
file->fp = fdopen(file->fd, "wb");
if (file->fp == NULL) {
fprintf(stderr, "fdopen memfile: %s\n", strerror(errno));
exit(EXIT_OPEN);
}
if (fwrite(file->map.c_str(), sizeof(char), file->map.size(), file->fp) != file->map.size()) {
fprintf(stderr, "memfile write: %s\n", strerror(errno));
exit(EXIT_WRITE);
}
}
+6 -7
View File
@@ -2,21 +2,20 @@
#define MEMFILE_HPP #define MEMFILE_HPP
#include <atomic> #include <atomic>
#include <string>
#include "errors.hpp"
struct memfile { struct memfile {
int fd = 0; int fd = 0;
char *map = NULL; std::string map;
std::atomic<long long> len;
long long off = 0;
unsigned long tree = 0; unsigned long tree = 0;
FILE *fp = NULL;
memfile() size_t off = 0;
: len(0) {
}
}; };
struct memfile *memfile_open(int fd); struct memfile *memfile_open(int fd);
int memfile_close(struct memfile *file); int memfile_close(struct memfile *file);
int memfile_write(struct memfile *file, void *s, long long len); int memfile_write(struct memfile *file, void *s, long long len);
void memfile_full(struct memfile *file);
#endif #endif
+2 -2
View File
@@ -82,14 +82,14 @@ int decompress(std::string const &input, std::string &output) {
} }
// https://github.com/mapbox/mapnik-vector-tile/blob/master/src/vector_tile_compression.hpp // https://github.com/mapbox/mapnik-vector-tile/blob/master/src/vector_tile_compression.hpp
int compress(std::string const &input, std::string &output) { int compress(std::string const &input, std::string &output, bool gz) {
z_stream deflate_s; z_stream deflate_s;
deflate_s.zalloc = Z_NULL; deflate_s.zalloc = Z_NULL;
deflate_s.zfree = Z_NULL; deflate_s.zfree = Z_NULL;
deflate_s.opaque = Z_NULL; deflate_s.opaque = Z_NULL;
deflate_s.avail_in = 0; deflate_s.avail_in = 0;
deflate_s.next_in = Z_NULL; deflate_s.next_in = Z_NULL;
deflateInit2(&deflate_s, Z_BEST_COMPRESSION, Z_DEFLATED, 31, 8, Z_DEFAULT_STRATEGY); deflateInit2(&deflate_s, Z_BEST_COMPRESSION, Z_DEFLATED, gz ? 31 : 15, 8, Z_DEFAULT_STRATEGY);
deflate_s.next_in = (Bytef *) input.data(); deflate_s.next_in = (Bytef *) input.data();
deflate_s.avail_in = input.size(); deflate_s.avail_in = input.size();
size_t length = 0; size_t length = 0;
+1 -1
View File
@@ -115,7 +115,7 @@ struct mvt_tile {
bool is_compressed(std::string const &data); bool is_compressed(std::string const &data);
int decompress(std::string const &input, std::string &output); int decompress(std::string const &input, std::string &output);
int compress(std::string const &input, std::string &output); int compress(std::string const &input, std::string &output, bool gz);
int dezig(unsigned n); int dezig(unsigned n);
mvt_value stringified_to_mvt_value(int type, const char *s); mvt_value stringified_to_mvt_value(int type, const char *s);
-1
View File
@@ -405,7 +405,6 @@ serial_feature parse_feature(json_pull *jp, int z, unsigned x, unsigned y, std::
sf.bbox[0] = sf.bbox[1] = LLONG_MAX; sf.bbox[0] = sf.bbox[1] = LLONG_MAX;
sf.bbox[2] = sf.bbox[3] = LLONG_MIN; sf.bbox[2] = sf.bbox[3] = LLONG_MIN;
sf.extent = 0; sf.extent = 0;
sf.metapos = 0;
sf.has_id = false; sf.has_id = false;
std::string layername = "unknown"; std::string layername = "unknown";
+3 -2
View File
@@ -57,7 +57,7 @@ std::string compress_fn(const std::string &input, uint8_t compression) {
if (compression == pmtiles::COMPRESSION_NONE) { if (compression == pmtiles::COMPRESSION_NONE) {
output = input; output = input;
} else if (compression == pmtiles::COMPRESSION_GZIP) { } else if (compression == pmtiles::COMPRESSION_GZIP) {
compress(input, output); compress(input, output, true);
} else { } else {
throw std::runtime_error("Unknown or unsupported compression."); throw std::runtime_error("Unknown or unsupported compression.");
} }
@@ -121,7 +121,7 @@ std::string metadata_to_pmtiles_json(metadata m) {
state.json_end_hash(); state.json_end_hash();
state.json_write_newline(); state.json_write_newline();
std::string compressed; std::string compressed;
compress(buf, compressed); compress(buf, compressed, true);
return compressed; return compressed;
} }
@@ -200,6 +200,7 @@ void mbtiles_map_image_to_pmtiles(char *fname, metadata m, bool tile_compression
pmtiles::zxy zxy = pmtiles::tileid_to_zxy(tile_id); pmtiles::zxy zxy = pmtiles::tileid_to_zxy(tile_id);
if (!quiet && !quiet_progress) { if (!quiet && !quiet_progress) {
fprintf(stderr, " %3.1f%% %d/%u/%u \r", progress, zxy.z, zxy.x, zxy.y); fprintf(stderr, " %3.1f%% %d/%u/%u \r", progress, zxy.z, zxy.x, zxy.y);
fflush(stderr);
} }
sqlite3_bind_int(map_stmt, 1, zxy.z); sqlite3_bind_int(map_stmt, 1, zxy.z);
sqlite3_bind_int(map_stmt, 2, zxy.x); sqlite3_bind_int(map_stmt, 2, zxy.x);
+39 -7
View File
@@ -3,6 +3,7 @@
#include <string.h> #include <string.h>
#include <limits.h> #include <limits.h>
#include <math.h> #include <math.h>
#include "main.hpp"
#include "memfile.hpp" #include "memfile.hpp"
#include "pool.hpp" #include "pool.hpp"
#include "errors.hpp" #include "errors.hpp"
@@ -42,24 +43,26 @@ long long addpool(struct memfile *poolfile, struct memfile *treefile, const char
} }
while (*sp != 0) { while (*sp != 0) {
int cmp = swizzlecmp(s, poolfile->map + ((struct stringpool *) (treefile->map + *sp))->off + 1); int cmp = swizzlecmp(s, poolfile->map.c_str() + ((struct stringpool *) (treefile->map.c_str() + *sp))->off + 1);
if (cmp == 0) { if (cmp == 0) {
cmp = type - (poolfile->map + ((struct stringpool *) (treefile->map + *sp))->off)[0]; cmp = type - (poolfile->map.c_str() + ((struct stringpool *) (treefile->map.c_str() + *sp))->off)[0];
} }
if (cmp < 0) { if (cmp < 0) {
sp = &(((struct stringpool *) (treefile->map + *sp))->left); sp = &(((struct stringpool *) (treefile->map.c_str() + *sp))->left);
} else if (cmp > 0) { } else if (cmp > 0) {
sp = &(((struct stringpool *) (treefile->map + *sp))->right); sp = &(((struct stringpool *) (treefile->map.c_str() + *sp))->right);
} else { } else {
return ((struct stringpool *) (treefile->map + *sp))->off; return ((struct stringpool *) (treefile->map.c_str() + *sp))->off;
} }
depth++; depth++;
if (depth > max) { if (depth > max) {
// Search is very deep, so string is probably unique. // Search is very deep, so string is probably unique.
// Add it to the pool without adding it to the search tree. // Add it to the pool without adding it to the search tree.
// This might go either to memory or the file, depending on whether
// the pool is full yet.
long long off = poolfile->off; long long off = poolfile->off;
if (memfile_write(poolfile, &type, 1) < 0) { if (memfile_write(poolfile, &type, 1) < 0) {
@@ -74,12 +77,41 @@ long long addpool(struct memfile *poolfile, struct memfile *treefile, const char
} }
} }
// Size of memory divided by 10 from observation of OOM errors (when supposedly
// 20% of memory is full) and onset of thrashing (when supposedly 15% of memory
// is full) on ECS.
if ((size_t) (poolfile->off + treefile->off) > memsize / CPUS / 10) {
// If the pool and search tree get to be larger than physical memory,
// then searching will start thrashing. Switch to appending strings
// to the file instead of keeping them in memory.
if (poolfile->fp == NULL) {
memfile_full(poolfile);
}
}
if (poolfile->fp != NULL) {
// We are now appending to the file, so don't try to keep tree references
// to the newly-added strings.
long long off = poolfile->off;
if (memfile_write(poolfile, &type, 1) < 0) {
perror("memfile write");
exit(EXIT_WRITE);
}
if (memfile_write(poolfile, (void *) s, strlen(s) + 1) < 0) {
perror("memfile write");
exit(EXIT_WRITE);
}
return off;
}
// *sp is probably in the memory-mapped file, and will move if the file grows. // *sp is probably in the memory-mapped file, and will move if the file grows.
long long ssp; long long ssp;
if (sp == &treefile->tree) { if (sp == &treefile->tree) {
ssp = -1; ssp = -1;
} else { } else {
ssp = ((char *) sp) - treefile->map; ssp = ((char *) sp) - treefile->map.c_str();
} }
long long off = poolfile->off; long long off = poolfile->off;
@@ -116,7 +148,7 @@ long long addpool(struct memfile *poolfile, struct memfile *treefile, const char
if (ssp == -1) { if (ssp == -1) {
treefile->tree = p; treefile->tree = p;
} else { } else {
*((long long *) (treefile->map + ssp)) = p; *((long long *) (treefile->map.c_str() + ssp)) = p;
} }
return off; return off;
} }
+108 -166
View File
@@ -9,9 +9,11 @@
#include <map> #include <map>
#include <algorithm> #include <algorithm>
#include <limits.h> #include <limits.h>
#include <zlib.h>
#include "protozero/varint.hpp" #include "protozero/varint.hpp"
#include "geometry.hpp" #include "geometry.hpp"
#include "mbtiles.hpp" #include "mbtiles.hpp"
#include "mvt.hpp"
#include "tile.hpp" #include "tile.hpp"
#include "serial.hpp" #include "serial.hpp"
#include "options.hpp" #include "options.hpp"
@@ -27,12 +29,15 @@
#define SHIFT_RIGHT(a) ((long long) std::round((double) (a) / (1LL << geometry_scale))) #define SHIFT_RIGHT(a) ((long long) std::round((double) (a) / (1LL << geometry_scale)))
#define SHIFT_LEFT(a) ((((a) + (COORD_OFFSET >> geometry_scale)) << geometry_scale) - COORD_OFFSET) #define SHIFT_LEFT(a) ((((a) + (COORD_OFFSET >> geometry_scale)) << geometry_scale) - COORD_OFFSET)
size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, const char *fname) { // write to file
size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, std::atomic<long long> *fpos, const char *fname) {
size_t w = fwrite(ptr, size, nitems, stream); size_t w = fwrite(ptr, size, nitems, stream);
if (w != nitems) { if (w != nitems) {
fprintf(stderr, "%s: Write to temporary file failed: %s\n", fname, strerror(errno)); fprintf(stderr, "%s: Write to temporary file failed: %s\n", fname, strerror(errno));
exit(EXIT_WRITE); exit(EXIT_WRITE);
} }
*fpos += size * nitems;
return w; return w;
} }
@@ -69,15 +74,54 @@ void serialize_ulong_long(FILE *out, unsigned long long zigzag, std::atomic<long
} }
void serialize_byte(FILE *out, signed char n, std::atomic<long long> *fpos, const char *fname) { void serialize_byte(FILE *out, signed char n, std::atomic<long long> *fpos, const char *fname) {
fwrite_check(&n, sizeof(signed char), 1, out, fname); fwrite_check(&n, sizeof(signed char), 1, out, fpos, fname);
*fpos += sizeof(signed char);
} }
void serialize_uint(FILE *out, unsigned n, std::atomic<long long> *fpos, const char *fname) { void serialize_uint(FILE *out, unsigned n, std::atomic<long long> *fpos, const char *fname) {
fwrite_check(&n, sizeof(unsigned), 1, out, fname); serialize_ulong_long(out, n, fpos, fname);
*fpos += sizeof(unsigned);
} }
// write to memory
size_t fwrite_check(const void *ptr, size_t size, size_t nitems, std::string &stream) {
stream += std::string((char *) ptr, size * nitems);
return nitems;
}
void serialize_ulong_long(std::string &out, unsigned long long zigzag) {
while (1) {
unsigned char b = zigzag & 0x7F;
if ((zigzag >> 7) != 0) {
b |= 0x80;
out += b;
zigzag >>= 7;
} else {
out += b;
break;
}
}
}
void serialize_long_long(std::string &out, long long n) {
unsigned long long zigzag = protozero::encode_zigzag64(n);
serialize_ulong_long(out, zigzag);
}
void serialize_int(std::string &out, int n) {
serialize_long_long(out, n);
}
void serialize_byte(std::string &out, signed char n) {
out += n;
}
void serialize_uint(std::string &out, unsigned n) {
serialize_ulong_long(out, n);
}
// read from memory
void deserialize_int(char **f, int *n) { void deserialize_int(char **f, int *n) {
long long ll; long long ll;
deserialize_long_long(f, &ll); deserialize_long_long(f, &ll);
@@ -109,8 +153,9 @@ void deserialize_ulong_long(char **f, unsigned long long *zigzag) {
} }
void deserialize_uint(char **f, unsigned *n) { void deserialize_uint(char **f, unsigned *n) {
memcpy(n, *f, sizeof(unsigned)); unsigned long long v;
*f += sizeof(unsigned); deserialize_ulong_long(f, &v);
*n = v;
} }
void deserialize_byte(char **f, signed char *n) { void deserialize_byte(char **f, signed char *n) {
@@ -118,79 +163,26 @@ void deserialize_byte(char **f, signed char *n) {
*f += sizeof(signed char); *f += sizeof(signed char);
} }
int deserialize_long_long_io(FILE *f, long long *n, std::atomic<long long> *geompos) { static void write_geometry(drawvec const &dv, std::string &out, long long wx, long long wy) {
unsigned long long zigzag = 0;
int ret = deserialize_ulong_long_io(f, &zigzag, geompos);
*n = protozero::decode_zigzag64(zigzag);
return ret;
}
int deserialize_ulong_long_io(FILE *f, unsigned long long *zigzag, std::atomic<long long> *geompos) {
*zigzag = 0;
int shift = 0;
while (1) {
int c = getc(f);
if (c == EOF) {
return 0;
}
(*geompos)++;
if ((c & 0x80) == 0) {
*zigzag |= ((unsigned long long) c) << shift;
shift += 7;
break;
} else {
*zigzag |= ((unsigned long long) (c & 0x7F)) << shift;
shift += 7;
}
}
return 1;
}
int deserialize_int_io(FILE *f, int *n, std::atomic<long long> *geompos) {
long long ll = 0;
int ret = deserialize_long_long_io(f, &ll, geompos);
*n = ll;
return ret;
}
int deserialize_uint_io(FILE *f, unsigned *n, std::atomic<long long> *geompos) {
if (fread(n, sizeof(unsigned), 1, f) != 1) {
return 0;
}
*geompos += sizeof(unsigned);
return 1;
}
int deserialize_byte_io(FILE *f, signed char *n, std::atomic<long long> *geompos) {
int c = getc(f);
if (c == EOF) {
return 0;
}
*n = c;
(*geompos)++;
return 1;
}
static void write_geometry(drawvec const &dv, std::atomic<long long> *fpos, FILE *out, const char *fname, long long wx, long long wy) {
for (size_t i = 0; i < dv.size(); i++) { for (size_t i = 0; i < dv.size(); i++) {
if (dv[i].op == VT_MOVETO || dv[i].op == VT_LINETO) { if (dv[i].op == VT_MOVETO || dv[i].op == VT_LINETO) {
serialize_byte(out, dv[i].op, fpos, fname); serialize_byte(out, dv[i].op);
serialize_long_long(out, dv[i].x - wx, fpos, fname); serialize_long_long(out, dv[i].x - wx);
serialize_long_long(out, dv[i].y - wy, fpos, fname); serialize_long_long(out, dv[i].y - wy);
wx = dv[i].x; wx = dv[i].x;
wy = dv[i].y; wy = dv[i].y;
} else { } else {
serialize_byte(out, dv[i].op, fpos, fname); serialize_byte(out, dv[i].op);
} }
} }
serialize_byte(out, VT_END);
} }
// called from generating the next zoom level // called from generating the next zoom level
void serialize_feature(FILE *geomfile, serial_feature *sf, std::atomic<long long> *geompos, const char *fname, long long wx, long long wy, bool include_minzoom) { std::string serialize_feature(serial_feature *sf, long long wx, long long wy) {
serialize_byte(geomfile, sf->t, geompos, fname); std::string s;
serialize_byte(s, sf->t);
#define FLAG_LAYER 7 #define FLAG_LAYER 7
@@ -212,63 +204,56 @@ void serialize_feature(FILE *geomfile, serial_feature *sf, std::atomic<long long
layer |= sf->has_tippecanoe_minzoom << FLAG_MINZOOM; layer |= sf->has_tippecanoe_minzoom << FLAG_MINZOOM;
layer |= sf->has_tippecanoe_maxzoom << FLAG_MAXZOOM; layer |= sf->has_tippecanoe_maxzoom << FLAG_MAXZOOM;
serialize_long_long(geomfile, layer, geompos, fname); serialize_long_long(s, layer);
if (sf->seq != 0) { if (sf->seq != 0) {
serialize_long_long(geomfile, sf->seq, geompos, fname); serialize_long_long(s, sf->seq);
} }
if (sf->has_tippecanoe_minzoom) { if (sf->has_tippecanoe_minzoom) {
serialize_int(geomfile, sf->tippecanoe_minzoom, geompos, fname); serialize_int(s, sf->tippecanoe_minzoom);
} }
if (sf->has_tippecanoe_maxzoom) { if (sf->has_tippecanoe_maxzoom) {
serialize_int(geomfile, sf->tippecanoe_maxzoom, geompos, fname); serialize_int(s, sf->tippecanoe_maxzoom);
} }
if (sf->has_id) { if (sf->has_id) {
serialize_ulong_long(geomfile, sf->id, geompos, fname); serialize_ulong_long(s, sf->id);
} }
serialize_int(geomfile, sf->segment, geompos, fname); serialize_int(s, sf->segment);
write_geometry(sf->geometry, s, wx, wy);
write_geometry(sf->geometry, geompos, geomfile, fname, wx, wy);
serialize_byte(geomfile, VT_END, geompos, fname);
if (sf->index != 0) { if (sf->index != 0) {
serialize_ulong_long(geomfile, sf->index, geompos, fname); serialize_ulong_long(s, sf->index);
} }
if (sf->label_point != 0) { if (sf->label_point != 0) {
serialize_ulong_long(geomfile, sf->label_point, geompos, fname); serialize_ulong_long(s, sf->label_point);
} }
if (sf->extent != 0) { if (sf->extent != 0) {
serialize_long_long(geomfile, sf->extent, geompos, fname); serialize_long_long(s, sf->extent);
} }
serialize_long_long(geomfile, sf->metapos, geompos, fname); serialize_long_long(s, sf->keys.size());
if (sf->metapos < 0) {
serialize_long_long(geomfile, sf->keys.size(), geompos, fname);
for (size_t i = 0; i < sf->keys.size(); i++) { for (size_t i = 0; i < sf->keys.size(); i++) {
serialize_long_long(geomfile, sf->keys[i], geompos, fname); serialize_long_long(s, sf->keys[i]);
serialize_long_long(geomfile, sf->values[i], geompos, fname); serialize_long_long(s, sf->values[i]);
}
} }
if (include_minzoom) { // MAGIC: This knows that the feature minzoom is the last byte of the feature,
serialize_byte(geomfile, sf->feature_minzoom, geompos, fname); serialize_byte(s, sf->feature_minzoom);
} return s;
} }
serial_feature deserialize_feature(FILE *geoms, std::atomic<long long> *geompos_in, char *metabase, long long *meta_off, unsigned z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y) { serial_feature deserialize_feature(std::string &geoms, unsigned z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y) {
serial_feature sf; serial_feature sf;
char *cp = (char *) geoms.c_str();
deserialize_byte_io(geoms, &sf.t, geompos_in); deserialize_byte(&cp, &sf.t);
if (sf.t < 0) { deserialize_long_long(&cp, &sf.layer);
return sf;
}
deserialize_long_long_io(geoms, &sf.layer, geompos_in);
sf.seq = 0; sf.seq = 0;
if (sf.layer & (1 << FLAG_SEQ)) { if (sf.layer & (1 << FLAG_SEQ)) {
deserialize_long_long_io(geoms, &sf.seq, geompos_in); deserialize_long_long(&cp, &sf.seq);
} }
sf.tippecanoe_minzoom = -1; sf.tippecanoe_minzoom = -1;
@@ -276,64 +261,54 @@ serial_feature deserialize_feature(FILE *geoms, std::atomic<long long> *geompos_
sf.id = 0; sf.id = 0;
sf.has_id = false; sf.has_id = false;
if (sf.layer & (1 << FLAG_MINZOOM)) { if (sf.layer & (1 << FLAG_MINZOOM)) {
deserialize_int_io(geoms, &sf.tippecanoe_minzoom, geompos_in); deserialize_int(&cp, &sf.tippecanoe_minzoom);
} }
if (sf.layer & (1 << FLAG_MAXZOOM)) { if (sf.layer & (1 << FLAG_MAXZOOM)) {
deserialize_int_io(geoms, &sf.tippecanoe_maxzoom, geompos_in); deserialize_int(&cp, &sf.tippecanoe_maxzoom);
} }
if (sf.layer & (1 << FLAG_ID)) { if (sf.layer & (1 << FLAG_ID)) {
sf.has_id = true; sf.has_id = true;
deserialize_ulong_long_io(geoms, &sf.id, geompos_in); deserialize_ulong_long(&cp, &sf.id);
} }
deserialize_int_io(geoms, &sf.segment, geompos_in); deserialize_int(&cp, &sf.segment);
sf.index = 0; sf.index = 0;
sf.label_point = 0; sf.label_point = 0;
sf.extent = 0; sf.extent = 0;
sf.geometry = decode_geometry(geoms, geompos_in, z, tx, ty, sf.bbox, initial_x[sf.segment], initial_y[sf.segment]); sf.geometry = decode_geometry(&cp, z, tx, ty, sf.bbox, initial_x[sf.segment], initial_y[sf.segment]);
if (sf.layer & (1 << FLAG_INDEX)) { if (sf.layer & (1 << FLAG_INDEX)) {
deserialize_ulong_long_io(geoms, &sf.index, geompos_in); deserialize_ulong_long(&cp, &sf.index);
} }
if (sf.layer & (1 << FLAG_LABEL_POINT)) { if (sf.layer & (1 << FLAG_LABEL_POINT)) {
deserialize_ulong_long_io(geoms, &sf.label_point, geompos_in); deserialize_ulong_long(&cp, &sf.label_point);
} }
if (sf.layer & (1 << FLAG_EXTENT)) { if (sf.layer & (1 << FLAG_EXTENT)) {
deserialize_long_long_io(geoms, &sf.extent, geompos_in); deserialize_long_long(&cp, &sf.extent);
} }
sf.layer >>= FLAG_LAYER; sf.layer >>= FLAG_LAYER;
sf.metapos = 0;
deserialize_long_long_io(geoms, &sf.metapos, geompos_in);
if (sf.metapos >= 0) {
char *meta = metabase + sf.metapos + meta_off[sf.segment];
long long count; long long count;
deserialize_long_long(&meta, &count); deserialize_long_long(&cp, &count);
for (long long i = 0; i < count; i++) { for (long long i = 0; i < count; i++) {
long long k, v; long long k, v;
deserialize_long_long(&meta, &k); deserialize_long_long(&cp, &k);
deserialize_long_long(&meta, &v); deserialize_long_long(&cp, &v);
sf.keys.push_back(k); sf.keys.push_back(k);
sf.values.push_back(v); sf.values.push_back(v);
} }
} else {
long long count;
deserialize_long_long_io(geoms, &count, geompos_in);
for (long long i = 0; i < count; i++) { // MAGIC: This knows that the feature minzoom is the last byte of the feature.
long long k, v; deserialize_byte(&cp, &sf.feature_minzoom);
deserialize_long_long_io(geoms, &k, geompos_in);
deserialize_long_long_io(geoms, &v, geompos_in);
sf.keys.push_back(k);
sf.values.push_back(v);
}
}
deserialize_byte_io(geoms, &sf.feature_minzoom, geompos_in); if (cp != geoms.c_str() + geoms.size()) {
fprintf(stderr, "wrong length decoding feature: used %zd, len is %zu\n", cp - geoms.c_str(), geoms.size());
exit(EXIT_IMPOSSIBLE);
}
return sf; return sf;
} }
@@ -512,32 +487,6 @@ int serialize_feature(struct serialization_state *sst, serial_feature &sf) {
locs.clear(); locs.clear();
} }
bool inline_meta = true;
// Don't inline metadata for features that will span several tiles at maxzoom
if (scaled_geometry.size() > 0 && (sf.bbox[2] < sf.bbox[0] || sf.bbox[3] < sf.bbox[1])) {
fprintf(stderr, "Internal error: impossible feature bounding box %llx,%llx,%llx,%llx\n", sf.bbox[0], sf.bbox[1], sf.bbox[2], sf.bbox[3]);
}
if (sf.bbox[0] == LLONG_MAX) {
// No bounding box (empty geometry)
// Shouldn't happen, but avoid arithmetic overflow below
} else if (sf.bbox[2] - sf.bbox[0] > (2LL << (32 - sst->maxzoom)) || sf.bbox[3] - sf.bbox[1] > (2LL << (32 - sst->maxzoom))) {
inline_meta = false;
if (prevent[P_CLIPPING]) {
static std::atomic<long long> warned(0);
long long extent = ((sf.bbox[2] - sf.bbox[0]) / ((1LL << (32 - sst->maxzoom)) + 1)) * ((sf.bbox[3] - sf.bbox[1]) / ((1LL << (32 - sst->maxzoom)) + 1));
if (extent > warned) {
fprintf(stderr, "Warning: %s:%d: Large unclipped (-pc) feature may be duplicated across %lld tiles\n", sst->fname, sst->line, extent);
warned = extent;
if (extent > 10000) {
fprintf(stderr, "Exiting because this can't be right.\n");
exit(EXIT_IMPOSSIBLE);
}
}
}
}
double extent = 0; double extent = 0;
if (additional[A_DROP_SMALLEST_AS_NEEDED] || additional[A_COALESCE_SMALLEST_AS_NEEDED] || order_by_size || sst->want_dist) { if (additional[A_DROP_SMALLEST_AS_NEEDED] || additional[A_COALESCE_SMALLEST_AS_NEEDED] || order_by_size || sst->want_dist) {
if (sf.t == VT_POLYGON) { if (sf.t == VT_POLYGON) {
@@ -725,24 +674,17 @@ int serialize_feature(struct serialization_state *sst, serial_feature &sf) {
} }
} }
if (inline_meta) {
sf.metapos = -1;
for (size_t i = 0; i < sf.full_keys.size(); i++) { for (size_t i = 0; i < sf.full_keys.size(); i++) {
sf.keys.push_back(addpool(r->poolfile, r->treefile, sf.full_keys[i].c_str(), mvt_string)); sf.keys.push_back(addpool(r->poolfile, r->treefile, sf.full_keys[i].c_str(), mvt_string));
sf.values.push_back(addpool(r->poolfile, r->treefile, sf.full_values[i].s.c_str(), sf.full_values[i].type)); sf.values.push_back(addpool(r->poolfile, r->treefile, sf.full_values[i].s.c_str(), sf.full_values[i].type));
} }
} else {
sf.metapos = r->metapos;
serialize_long_long(r->metafile, sf.full_keys.size(), &r->metapos, sst->fname);
for (size_t i = 0; i < sf.full_keys.size(); i++) {
serialize_long_long(r->metafile, addpool(r->poolfile, r->treefile, sf.full_keys[i].c_str(), mvt_string), &r->metapos, sst->fname);
serialize_long_long(r->metafile, addpool(r->poolfile, r->treefile, sf.full_values[i].s.c_str(), sf.full_values[i].type), &r->metapos, sst->fname);
}
}
long long geomstart = r->geompos; long long geomstart = r->geompos;
sf.geometry = scaled_geometry; sf.geometry = scaled_geometry;
serialize_feature(r->geomfile, &sf, &r->geompos, sst->fname, SHIFT_RIGHT(*(sst->initial_x)), SHIFT_RIGHT(*(sst->initial_y)), false);
std::string feature = serialize_feature(&sf, SHIFT_RIGHT(*(sst->initial_x)), SHIFT_RIGHT(*(sst->initial_y)));
serialize_long_long(r->geomfile, feature.size(), &r->geompos, sst->fname);
fwrite_check(feature.c_str(), sizeof(char), feature.size(), r->geomfile, &r->geompos, sst->fname);
struct index index; struct index index;
index.start = geomstart; index.start = geomstart;
@@ -752,8 +694,7 @@ int serialize_feature(struct serialization_state *sst, serial_feature &sf) {
index.t = sf.t; index.t = sf.t;
index.ix = bbox_index; index.ix = bbox_index;
fwrite_check(&index, sizeof(struct index), 1, r->indexfile, sst->fname); fwrite_check(&index, sizeof(struct index), 1, r->indexfile, &r->indexpos, sst->fname);
r->indexpos += sizeof(struct index);
for (size_t i = 0; i < 2; i++) { for (size_t i = 0; i < 2; i++) {
if (sf.bbox[i] < r->file_bbox[i]) { if (sf.bbox[i] < r->file_bbox[i]) {
@@ -770,6 +711,7 @@ int serialize_feature(struct serialization_state *sst, serial_feature &sf) {
checkdisk(sst->readers); checkdisk(sst->readers);
if (!quiet && !quiet_progress && progress_time()) { if (!quiet && !quiet_progress && progress_time()) {
fprintf(stderr, "Read %.2f million features\r", *sst->progress_seq / 1000000.0); fprintf(stderr, "Read %.2f million features\r", *sst->progress_seq / 1000000.0);
fflush(stderr);
} }
} }
(*(sst->progress_seq))++; (*(sst->progress_seq))++;
+11 -24
View File
@@ -11,14 +11,19 @@
#include "mbtiles.hpp" #include "mbtiles.hpp"
#include "jsonpull/jsonpull.h" #include "jsonpull/jsonpull.h"
size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, const char *fname); size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, std::atomic<long long> *fpos, const char *fname);
void serialize_int(FILE *out, int n, std::atomic<long long> *fpos, const char *fname); void serialize_int(FILE *out, int n, std::atomic<long long> *fpos, const char *fname);
void serialize_long_long(FILE *out, long long n, std::atomic<long long> *fpos, const char *fname); void serialize_long_long(FILE *out, long long n, std::atomic<long long> *fpos, const char *fname);
void serialize_ulong_long(FILE *out, unsigned long long n, std::atomic<long long> *fpos, const char *fname); void serialize_ulong_long(FILE *out, unsigned long long n, std::atomic<long long> *fpos, const char *fname);
void serialize_byte(FILE *out, signed char n, std::atomic<long long> *fpos, const char *fname); void serialize_byte(FILE *out, signed char n, std::atomic<long long> *fpos, const char *fname);
void serialize_uint(FILE *out, unsigned n, std::atomic<long long> *fpos, const char *fname); void serialize_uint(FILE *out, unsigned n, std::atomic<long long> *fpos, const char *fname);
void serialize_string(FILE *out, const char *s, std::atomic<long long> *fpos, const char *fname);
void serialize_int(std::string &out, int n);
void serialize_long_long(std::string &out, long long n);
void serialize_ulong_long(std::string &out, unsigned long long n);
void serialize_byte(std::string &out, signed char n);
void serialize_uint(std::string &out, unsigned n);
void deserialize_int(char **f, int *n); void deserialize_int(char **f, int *n);
void deserialize_long_long(char **f, long long *n); void deserialize_long_long(char **f, long long *n);
@@ -26,12 +31,6 @@ void deserialize_ulong_long(char **f, unsigned long long *n);
void deserialize_uint(char **f, unsigned *n); void deserialize_uint(char **f, unsigned *n);
void deserialize_byte(char **f, signed char *n); void deserialize_byte(char **f, signed char *n);
int deserialize_int_io(FILE *f, int *n, std::atomic<long long> *geompos);
int deserialize_long_long_io(FILE *f, long long *n, std::atomic<long long> *geompos);
int deserialize_ulong_long_io(FILE *f, unsigned long long *n, std::atomic<long long> *geompos);
int deserialize_uint_io(FILE *f, unsigned *n, std::atomic<long long> *geompos);
int deserialize_byte_io(FILE *f, signed char *n, std::atomic<long long> *geompos);
struct serial_val { struct serial_val {
int type = 0; int type = 0;
std::string s = ""; std::string s = "";
@@ -61,8 +60,6 @@ struct serial_feature {
std::vector<long long> keys{}; std::vector<long long> keys{};
std::vector<long long> values{}; std::vector<long long> values{};
// If >= 0, metadata is external
long long metapos = 0;
// XXX This isn't serialized. Should it be here? // XXX This isn't serialized. Should it be here?
long long bbox[4] = {0, 0, 0, 0}; long long bbox[4] = {0, 0, 0, 0};
@@ -72,54 +69,45 @@ struct serial_feature {
bool dropped = false; bool dropped = false;
}; };
void serialize_feature(FILE *geomfile, serial_feature *sf, std::atomic<long long> *geompos, const char *fname, long long wx, long long wy, bool include_minzoom); std::string serialize_feature(serial_feature *sf, long long wx, long long wy);
serial_feature deserialize_feature(FILE *geoms, std::atomic<long long> *geompos_in, char *metabase, long long *meta_off, unsigned z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y); serial_feature deserialize_feature(std::string &geoms, unsigned z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y);
struct reader { struct reader {
int metafd = -1;
int poolfd = -1; int poolfd = -1;
int treefd = -1; int treefd = -1;
int geomfd = -1; int geomfd = -1;
int indexfd = -1; int indexfd = -1;
FILE *metafile = NULL;
struct memfile *poolfile = NULL; struct memfile *poolfile = NULL;
struct memfile *treefile = NULL; struct memfile *treefile = NULL;
FILE *geomfile = NULL; FILE *geomfile = NULL;
FILE *indexfile = NULL; FILE *indexfile = NULL;
std::atomic<long long> metapos;
std::atomic<long long> geompos; std::atomic<long long> geompos;
std::atomic<long long> indexpos; std::atomic<long long> indexpos;
long long file_bbox[4] = {0, 0, 0, 0}; long long file_bbox[4] = {0, 0, 0, 0};
struct stat geomst {}; struct stat geomst {};
struct stat metast {};
char *geom_map = NULL; char *geom_map = NULL;
reader() reader()
: metapos(0), geompos(0), indexpos(0) { : geompos(0), indexpos(0) {
} }
reader(reader const &r) { reader(reader const &r) {
metafd = r.metafd;
poolfd = r.poolfd; poolfd = r.poolfd;
treefd = r.treefd; treefd = r.treefd;
geomfd = r.geomfd; geomfd = r.geomfd;
indexfd = r.indexfd; indexfd = r.indexfd;
metafile = r.metafile;
poolfile = r.poolfile; poolfile = r.poolfile;
treefile = r.treefile; treefile = r.treefile;
geomfile = r.geomfile; geomfile = r.geomfile;
indexfile = r.indexfile; indexfile = r.indexfile;
long long p = r.metapos; long long p = r.geompos;
metapos = p;
p = r.geompos;
geompos = p; geompos = p;
p = r.indexpos; p = r.indexpos;
@@ -128,7 +116,6 @@ struct reader {
memcpy(file_bbox, r.file_bbox, sizeof(file_bbox)); memcpy(file_bbox, r.file_bbox, sizeof(file_bbox));
geomst = r.geomst; geomst = r.geomst;
metast = r.metast;
geom_map = r.geom_map; geom_map = r.geom_map;
} }
File diff suppressed because one or more lines are too long
Binary file not shown.
+2 -2
View File
@@ -25,7 +25,6 @@
#include <pthread.h> #include <pthread.h>
#include "mvt.hpp" #include "mvt.hpp"
#include "projection.hpp" #include "projection.hpp"
#include "pool.hpp"
#include "mbtiles.hpp" #include "mbtiles.hpp"
#include "geometry.hpp" #include "geometry.hpp"
#include "dirtiles.hpp" #include "dirtiles.hpp"
@@ -545,7 +544,7 @@ void *join_worker(void *v) {
std::string compressed; std::string compressed;
if (!pC) { if (!pC) {
compress(pbf, compressed); compress(pbf, compressed, true);
} else { } else {
compressed = pbf; compressed = pbf;
} }
@@ -589,6 +588,7 @@ void handle_tasks(std::map<zxy, std::vector<std::string>> &tasks, std::vector<st
if (ai == tasks.begin()) { if (ai == tasks.begin()) {
if (!quiet) { if (!quiet) {
fprintf(stderr, "%lld/%lld/%lld \r", ai->first.z, ai->first.x, ai->first.y); fprintf(stderr, "%lld/%lld/%lld \r", ai->first.z, ai->first.x, ai->first.y);
fflush(stderr);
} }
} }
} }
+104 -41
View File
@@ -25,6 +25,7 @@
#include <errno.h> #include <errno.h>
#include <time.h> #include <time.h>
#include <fcntl.h> #include <fcntl.h>
#include <zlib.h>
#include <sys/wait.h> #include <sys/wait.h>
#include "mvt.hpp" #include "mvt.hpp"
#include "mbtiles.hpp" #include "mbtiles.hpp"
@@ -40,6 +41,8 @@
#include "milo/dtoa_milo.h" #include "milo/dtoa_milo.h"
#include "evaluator.hpp" #include "evaluator.hpp"
#include "errors.hpp" #include "errors.hpp"
#include "compression.hpp"
#include "protozero/varint.hpp"
extern "C" { extern "C" {
#include "jsonpull/jsonpull.h" #include "jsonpull/jsonpull.h"
@@ -328,7 +331,7 @@ struct ordercmp {
} }
} ordercmp; } ordercmp;
void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, unsigned tx, unsigned ty, int buffer, int *within, std::atomic<long long> *geompos, FILE **geomfile, const char *fname, signed char t, int layer, long long metastart, signed char feature_minzoom, int child_shards, int max_zoom_increment, long long seq, int tippecanoe_minzoom, int tippecanoe_maxzoom, int segment, unsigned *initial_x, unsigned *initial_y, std::vector<long long> &metakeys, std::vector<long long> &metavals, bool has_id, unsigned long long id, unsigned long long index, unsigned long long label_point, long long extent) { void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, unsigned tx, unsigned ty, int buffer, int *within, std::atomic<long long> *geompos, compressor **geomfile, const char *fname, signed char t, int layer, signed char feature_minzoom, int child_shards, int max_zoom_increment, long long seq, int tippecanoe_minzoom, int tippecanoe_maxzoom, int segment, unsigned *initial_x, unsigned *initial_y, std::vector<long long> &metakeys, std::vector<long long> &metavals, bool has_id, unsigned long long id, unsigned long long index, unsigned long long label_point, long long extent) {
if (geom.size() > 0 && (nextzoom <= maxzoom || additional[A_EXTEND_ZOOMS])) { if (geom.size() > 0 && (nextzoom <= maxzoom || additional[A_EXTEND_ZOOMS])) {
int xo, yo; int xo, yo;
int span = 1 << (nextzoom - z); int span = 1 << (nextzoom - z);
@@ -398,9 +401,10 @@ void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, u
{ {
if (!within[j]) { if (!within[j]) {
serialize_int(geomfile[j], nextzoom, &geompos[j], fname); serialize_int(geomfile[j]->fp, nextzoom, &geompos[j], fname);
serialize_uint(geomfile[j], tx * span + xo, &geompos[j], fname); serialize_uint(geomfile[j]->fp, tx * span + xo, &geompos[j], fname);
serialize_uint(geomfile[j], ty * span + yo, &geompos[j], fname); serialize_uint(geomfile[j]->fp, ty * span + yo, &geompos[j], fname);
geomfile[j]->begin();
within[j] = 1; within[j] = 1;
} }
@@ -415,21 +419,20 @@ void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, u
sf.tippecanoe_minzoom = tippecanoe_minzoom; sf.tippecanoe_minzoom = tippecanoe_minzoom;
sf.has_tippecanoe_maxzoom = tippecanoe_maxzoom != -1; sf.has_tippecanoe_maxzoom = tippecanoe_maxzoom != -1;
sf.tippecanoe_maxzoom = tippecanoe_maxzoom; sf.tippecanoe_maxzoom = tippecanoe_maxzoom;
sf.metapos = metastart;
sf.geometry = geom2; sf.geometry = geom2;
sf.index = index; sf.index = index;
sf.label_point = label_point; sf.label_point = label_point;
sf.extent = extent; sf.extent = extent;
sf.feature_minzoom = feature_minzoom; sf.feature_minzoom = feature_minzoom;
if (metastart < 0) {
for (size_t i = 0; i < metakeys.size(); i++) { for (size_t i = 0; i < metakeys.size(); i++) {
sf.keys.push_back(metakeys[i]); sf.keys.push_back(metakeys[i]);
sf.values.push_back(metavals[i]); sf.values.push_back(metavals[i]);
} }
}
serialize_feature(geomfile[j], &sf, &geompos[j], fname, SHIFT_RIGHT(initial_x[segment]), SHIFT_RIGHT(initial_y[segment]), true); std::string feature = serialize_feature(&sf, SHIFT_RIGHT(initial_x[segment]), SHIFT_RIGHT(initial_y[segment]));
geomfile[j]->serialize_long_long(feature.size(), &geompos[j], fname);
geomfile[j]->fwrite_check(feature.c_str(), sizeof(char), feature.size(), &geompos[j], fname);
} }
} }
} }
@@ -1287,14 +1290,13 @@ long long choose_minextent(std::vector<long long> &extents, double f) {
struct write_tile_args { struct write_tile_args {
struct task *tasks = NULL; struct task *tasks = NULL;
char *metabase = NULL;
char *stringpool = NULL; char *stringpool = NULL;
int min_detail = 0; int min_detail = 0;
sqlite3 *outdb = NULL; sqlite3 *outdb = NULL;
const char *outdir = NULL; const char *outdir = NULL;
int buffer = 0; int buffer = 0;
const char *fname = NULL; const char *fname = NULL;
FILE **geomfile = NULL; compressor **geomfile = NULL;
double todo = 0; double todo = 0;
std::atomic<long long> *along = NULL; std::atomic<long long> *along = NULL;
double gamma = 0; double gamma = 0;
@@ -1310,7 +1312,6 @@ struct write_tile_args {
int low_detail = 0; int low_detail = 0;
double simplification = 0; double simplification = 0;
std::atomic<long long> *most = NULL; std::atomic<long long> *most = NULL;
long long *meta_off = NULL;
long long *pool_off = NULL; long long *pool_off = NULL;
unsigned *initial_x = NULL; unsigned *initial_x = NULL;
unsigned *initial_y = NULL; unsigned *initial_y = NULL;
@@ -1336,6 +1337,8 @@ struct write_tile_args {
struct json_object *filter = NULL; struct json_object *filter = NULL;
std::atomic<size_t> *dropped_count = NULL; std::atomic<size_t> *dropped_count = NULL;
atomic_strategy *strategy = NULL; atomic_strategy *strategy = NULL;
int zoom = -1;
bool compressed;
}; };
bool clip_to_tile(serial_feature &sf, int z, long long buffer) { bool clip_to_tile(serial_feature &sf, int z, long long buffer) {
@@ -1433,18 +1436,40 @@ void remove_attributes(serial_feature &sf, std::set<std::string> const &exclude_
} }
} }
serial_feature next_feature(FILE *geoms, std::atomic<long long> *geompos_in, char *metabase, long long *meta_off, int z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y, long long *original_features, long long *unclipped_features, int nextzoom, int maxzoom, int minzoom, int max_zoom_increment, size_t pass, std::atomic<long long> *along, long long alongminus, int buffer, int *within, FILE **geomfile, std::atomic<long long> *geompos, std::atomic<double> *oprogress, double todo, const char *fname, int child_shards, struct json_object *filter, const char *stringpool, long long *pool_off, std::vector<std::vector<std::string>> *layer_unmaps, bool first_time) { serial_feature next_feature(decompressor *geoms, std::atomic<long long> *geompos_in, int z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y, long long *original_features, long long *unclipped_features, int nextzoom, int maxzoom, int minzoom, int max_zoom_increment, size_t pass, std::atomic<long long> *along, long long alongminus, int buffer, int *within, compressor **geomfile, std::atomic<long long> *geompos, std::atomic<double> *oprogress, double todo, const char *fname, int child_shards, struct json_object *filter, const char *stringpool, long long *pool_off, std::vector<std::vector<std::string>> *layer_unmaps, bool first_time, bool compressed) {
while (1) { while (1) {
serial_feature sf = deserialize_feature(geoms, geompos_in, metabase, meta_off, z, tx, ty, initial_x, initial_y); serial_feature sf;
if (sf.t < 0) { std::string s;
long long len;
if (geoms->deserialize_long_long(&len, geompos_in) == 0) {
fprintf(stderr, "Unexpected physical EOF in feature stream\n");
exit(EXIT_READ);
}
if (len == 0) {
if (compressed) {
geoms->end(geompos_in);
}
sf.t = -2;
return sf; return sf;
} }
s.resize(std::abs(len));
size_t n = geoms->fread((void *) s.c_str(), sizeof(char), s.size(), geompos_in);
if (n != s.size()) {
fprintf(stderr, "Short read (%zu for %zu) from geometry\n", n, s.size());
exit(EXIT_READ);
}
sf = deserialize_feature(s, z, tx, ty, initial_x, initial_y);
size_t passes = pass + 1; size_t passes = pass + 1;
double progress = floor(((((*geompos_in + *along - alongminus) / (double) todo) + pass) / passes + z) / (maxzoom + 1) * 1000) / 10; double progress = floor(((((*geompos_in + *along - alongminus) / (double) todo) + pass) / passes + z) / (maxzoom + 1) * 1000) / 10;
if (progress >= *oprogress + 0.1) { if (progress >= *oprogress + 0.1) {
if (!quiet && !quiet_progress && progress_time()) { if (!quiet && !quiet_progress && progress_time()) {
fprintf(stderr, " %3.1f%% %d/%u/%u \r", progress, z, tx, ty); fprintf(stderr, " %3.1f%% %d/%u/%u \r", progress, z, tx, ty);
fflush(stderr);
} }
if (logger.json_enabled && progress_time()) { if (logger.json_enabled && progress_time()) {
logger.progress_tile(progress); logger.progress_tile(progress);
@@ -1464,7 +1489,7 @@ serial_feature next_feature(FILE *geoms, std::atomic<long long> *geompos_in, cha
if (first_time && pass == 0) { /* only write out the next zoom once, even if we retry */ if (first_time && pass == 0) { /* only write out the next zoom once, even if we retry */
if (sf.tippecanoe_maxzoom == -1 || sf.tippecanoe_maxzoom >= nextzoom) { if (sf.tippecanoe_maxzoom == -1 || sf.tippecanoe_maxzoom >= nextzoom) {
rewrite(sf.geometry, z, nextzoom, maxzoom, sf.bbox, tx, ty, buffer, within, geompos, geomfile, fname, sf.t, sf.layer, sf.metapos, sf.feature_minzoom, child_shards, max_zoom_increment, sf.seq, sf.tippecanoe_minzoom, sf.tippecanoe_maxzoom, sf.segment, initial_x, initial_y, sf.keys, sf.values, sf.has_id, sf.id, sf.index, sf.label_point, sf.extent); rewrite(sf.geometry, z, nextzoom, maxzoom, sf.bbox, tx, ty, buffer, within, geompos, geomfile, fname, sf.t, sf.layer, sf.feature_minzoom, child_shards, max_zoom_increment, sf.seq, sf.tippecanoe_minzoom, sf.tippecanoe_maxzoom, sf.segment, initial_x, initial_y, sf.keys, sf.values, sf.has_id, sf.id, sf.index, sf.label_point, sf.extent);
} }
} }
@@ -1565,10 +1590,8 @@ serial_feature next_feature(FILE *geoms, std::atomic<long long> *geompos_in, cha
} }
struct run_prefilter_args { struct run_prefilter_args {
FILE *geoms = NULL; decompressor *geoms = NULL;
std::atomic<long long> *geompos_in = NULL; std::atomic<long long> *geompos_in = NULL;
char *metabase = NULL;
long long *meta_off = NULL;
int z = 0; int z = 0;
unsigned tx = 0; unsigned tx = 0;
unsigned ty = 0; unsigned ty = 0;
@@ -1585,7 +1608,7 @@ struct run_prefilter_args {
long long alongminus = 0; long long alongminus = 0;
int buffer = 0; int buffer = 0;
int *within = NULL; int *within = NULL;
FILE **geomfile = NULL; compressor **geomfile = NULL;
std::atomic<long long> *geompos = NULL; std::atomic<long long> *geompos = NULL;
std::atomic<double> *oprogress = NULL; std::atomic<double> *oprogress = NULL;
double todo = 0; double todo = 0;
@@ -1597,6 +1620,7 @@ struct run_prefilter_args {
FILE *prefilter_fp = NULL; FILE *prefilter_fp = NULL;
struct json_object *filter = NULL; struct json_object *filter = NULL;
bool first_time = false; bool first_time = false;
bool compressed = false;
}; };
void *run_prefilter(void *v) { void *run_prefilter(void *v) {
@@ -1604,7 +1628,7 @@ void *run_prefilter(void *v) {
json_writer state(rpa->prefilter_fp); json_writer state(rpa->prefilter_fp);
while (1) { while (1) {
serial_feature sf = next_feature(rpa->geoms, rpa->geompos_in, rpa->metabase, rpa->meta_off, rpa->z, rpa->tx, rpa->ty, rpa->initial_x, rpa->initial_y, rpa->original_features, rpa->unclipped_features, rpa->nextzoom, rpa->maxzoom, rpa->minzoom, rpa->max_zoom_increment, rpa->pass, rpa->along, rpa->alongminus, rpa->buffer, rpa->within, rpa->geomfile, rpa->geompos, rpa->oprogress, rpa->todo, rpa->fname, rpa->child_shards, rpa->filter, rpa->stringpool, rpa->pool_off, rpa->layer_unmaps, rpa->first_time); serial_feature sf = next_feature(rpa->geoms, rpa->geompos_in, rpa->z, rpa->tx, rpa->ty, rpa->initial_x, rpa->initial_y, rpa->original_features, rpa->unclipped_features, rpa->nextzoom, rpa->maxzoom, rpa->minzoom, rpa->max_zoom_increment, rpa->pass, rpa->along, rpa->alongminus, rpa->buffer, rpa->within, rpa->geomfile, rpa->geompos, rpa->oprogress, rpa->todo, rpa->fname, rpa->child_shards, rpa->filter, rpa->stringpool, rpa->pool_off, rpa->layer_unmaps, rpa->first_time, rpa->compressed);
if (sf.t < 0) { if (sf.t < 0) {
break; break;
} }
@@ -1854,7 +1878,7 @@ void add_sample_to(std::vector<T> &vals, T val, size_t &increment, size_t seq) {
} }
} }
long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *metabase, char *stringpool, int z, const unsigned tx, const unsigned ty, const int detail, int min_detail, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, FILE **geomfile, int minzoom, int maxzoom, double todo, std::atomic<long long> *along, long long alongminus, double gamma, int child_shards, long long *meta_off, long long *pool_off, unsigned *initial_x, unsigned *initial_y, std::atomic<int> *running, double simplification, std::vector<std::map<std::string, layermap_entry>> *layermaps, std::vector<std::vector<std::string>> *layer_unmaps, size_t tiling_seg, size_t pass, unsigned long long mingap, long long minextent, double fraction, const char *prefilter, const char *postfilter, struct json_object *filter, write_tile_args *arg, atomic_strategy *strategy) { long long write_tile(decompressor *geoms, std::atomic<long long> *geompos_in, char *stringpool, int z, const unsigned tx, const unsigned ty, const int detail, int min_detail, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, compressor **geomfile, int minzoom, int maxzoom, double todo, std::atomic<long long> *along, long long alongminus, double gamma, int child_shards, long long *pool_off, unsigned *initial_x, unsigned *initial_y, std::atomic<int> *running, double simplification, std::vector<std::map<std::string, layermap_entry>> *layermaps, std::vector<std::vector<std::string>> *layer_unmaps, size_t tiling_seg, size_t pass, unsigned long long mingap, long long minextent, double fraction, const char *prefilter, const char *postfilter, struct json_object *filter, write_tile_args *arg, atomic_strategy *strategy, bool compressed_input) {
double merge_fraction = 1; double merge_fraction = 1;
double mingap_fraction = 1; double mingap_fraction = 1;
double minextent_fraction = 1; double minextent_fraction = 1;
@@ -1926,11 +1950,22 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
} }
if (*geompos_in != og) { if (*geompos_in != og) {
if (fseek(geoms, og, SEEK_SET) != 0) { if (compressed_input) {
if (geoms->within) {
geoms->end(geompos_in);
}
geoms->begin();
}
if (fseek(geoms->fp, og, SEEK_SET) != 0) {
perror("fseek geom"); perror("fseek geom");
exit(EXIT_SEEK); exit(EXIT_SEEK);
} }
*geompos_in = og; *geompos_in = og;
geoms->zs.avail_in = 0;
geoms->zs.avail_out = 0;
} }
int prefilter_write = -1, prefilter_read = -1; int prefilter_write = -1, prefilter_read = -1;
@@ -1958,8 +1993,6 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
rpa.geoms = geoms; rpa.geoms = geoms;
rpa.geompos_in = geompos_in; rpa.geompos_in = geompos_in;
rpa.metabase = metabase;
rpa.meta_off = meta_off;
rpa.z = z; rpa.z = z;
rpa.tx = tx; rpa.tx = tx;
rpa.ty = ty; rpa.ty = ty;
@@ -1988,6 +2021,7 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
rpa.pool_off = pool_off; rpa.pool_off = pool_off;
rpa.filter = filter; rpa.filter = filter;
rpa.first_time = first_time; rpa.first_time = first_time;
rpa.compressed = compressed_input;
if (pthread_create(&prefilter_writer, NULL, run_prefilter, &rpa) != 0) { if (pthread_create(&prefilter_writer, NULL, run_prefilter, &rpa) != 0) {
perror("pthread_create (prefilter writer)"); perror("pthread_create (prefilter writer)");
@@ -2007,7 +2041,7 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
ssize_t which_partial = -1; ssize_t which_partial = -1;
if (prefilter == NULL) { if (prefilter == NULL) {
sf = next_feature(geoms, geompos_in, metabase, meta_off, z, tx, ty, initial_x, initial_y, &original_features, &unclipped_features, nextzoom, maxzoom, minzoom, max_zoom_increment, pass, along, alongminus, buffer, within, geomfile, geompos, &oprogress, todo, fname, child_shards, filter, stringpool, pool_off, layer_unmaps, first_time); sf = next_feature(geoms, geompos_in, z, tx, ty, initial_x, initial_y, &original_features, &unclipped_features, nextzoom, maxzoom, minzoom, max_zoom_increment, pass, along, alongminus, buffer, within, geomfile, geompos, &oprogress, todo, fname, child_shards, filter, stringpool, pool_off, layer_unmaps, first_time, compressed_input);
} else { } else {
sf = parse_feature(prefilter_jp, z, tx, ty, layermaps, tiling_seg, layer_unmaps, postfilter != NULL); sf = parse_feature(prefilter_jp, z, tx, ty, layermaps, tiling_seg, layer_unmaps, postfilter != NULL);
} }
@@ -2414,7 +2448,8 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
int j; int j;
for (j = 0; j < child_shards; j++) { for (j = 0; j < child_shards; j++) {
if (within[j]) { if (within[j]) {
serialize_byte(geomfile[j], -2, &geompos[j], fname); geomfile[j]->serialize_long_long(0, &geompos[j], fname); // EOF
geomfile[j]->end(&geompos[j], fname);
within[j] = 0; within[j] = 0;
} }
} }
@@ -2577,6 +2612,7 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
if (progress >= oprogress + 0.1) { if (progress >= oprogress + 0.1) {
if (!quiet && !quiet_progress && progress_time()) { if (!quiet && !quiet_progress && progress_time()) {
fprintf(stderr, " %3.1f%% %d/%u/%u \r", progress, z, tx, ty); fprintf(stderr, " %3.1f%% %d/%u/%u \r", progress, z, tx, ty);
fflush(stderr);
} }
if (logger.json_enabled && progress_time()) { if (logger.json_enabled && progress_time()) {
logger.progress_tile(progress); logger.progress_tile(progress);
@@ -2672,7 +2708,7 @@ long long write_tile(FILE *geoms, std::atomic<long long> *geompos_in, char *meta
std::string pbf = tile.encode(); std::string pbf = tile.encode();
if (!prevent[P_TILE_COMPRESSION]) { if (!prevent[P_TILE_COMPRESSION]) {
compress(pbf, compressed); compress(pbf, compressed, true);
} else { } else {
compressed = pbf; compressed = pbf;
} }
@@ -2818,7 +2854,15 @@ void *run_thread(void *vargs) {
continue; continue;
} }
// printf("%lld of geom_size\n", (long long) geom_size[j]); // If this is zoom level 0, the geomfd will be uncompressed data,
// because (at least for now) it needs to stay uncompressed during
// the sort and post-sort maxzoom calculation and fixup so that
// the sort can rearrange individual features and the fixup can
// then adjust their minzooms without decompressing and recompressing
// each feature.
//
// In higher zooms, it will be compressed data written out during the
// previous zoom.
FILE *geom = fdopen(arg->geomfd[j], "rb"); FILE *geom = fdopen(arg->geomfd[j], "rb");
if (geom == NULL) { if (geom == NULL) {
@@ -2826,6 +2870,8 @@ void *run_thread(void *vargs) {
exit(EXIT_OPEN); exit(EXIT_OPEN);
} }
decompressor dc(geom);
std::atomic<long long> geompos(0); std::atomic<long long> geompos(0);
long long prevgeom = 0; long long prevgeom = 0;
@@ -2833,17 +2879,31 @@ void *run_thread(void *vargs) {
int z; int z;
unsigned x, y; unsigned x, y;
if (!deserialize_int_io(geom, &z, &geompos)) { // These z/x/y are uncompressed so we can seek to the start of the
// compressed feature data that immediately follows.
if (!dc.deserialize_int(&z, &geompos)) {
break; break;
} }
deserialize_uint_io(geom, &x, &geompos); dc.deserialize_uint(&x, &geompos);
deserialize_uint_io(geom, &y, &geompos); dc.deserialize_uint(&y, &geompos);
#if 0
// currently broken because also requires tracking nextzoom when skipping zooms
if (z != arg->zoom) {
fprintf(stderr, "Expected zoom %d, found zoom %d\n", arg->zoom, z);
exit(EXIT_IMPOSSIBLE);
}
#endif
if (arg->compressed) {
dc.begin();
}
arg->wrote_zoom = z; arg->wrote_zoom = z;
// fprintf(stderr, "%d/%u/%u\n", z, x, y); // fprintf(stderr, "%d/%u/%u\n", z, x, y);
long long len = write_tile(geom, &geompos, arg->metabase, arg->stringpool, z, x, y, z == arg->maxzoom ? arg->full_detail : arg->low_detail, arg->min_detail, arg->outdb, arg->outdir, arg->buffer, arg->fname, arg->geomfile, arg->minzoom, arg->maxzoom, arg->todo, arg->along, geompos, arg->gamma, arg->child_shards, arg->meta_off, arg->pool_off, arg->initial_x, arg->initial_y, arg->running, arg->simplification, arg->layermaps, arg->layer_unmaps, arg->tiling_seg, arg->pass, arg->mingap, arg->minextent, arg->fraction, arg->prefilter, arg->postfilter, arg->filter, arg, arg->strategy); long long len = write_tile(&dc, &geompos, arg->stringpool, z, x, y, z == arg->maxzoom ? arg->full_detail : arg->low_detail, arg->min_detail, arg->outdb, arg->outdir, arg->buffer, arg->fname, arg->geomfile, arg->minzoom, arg->maxzoom, arg->todo, arg->along, geompos, arg->gamma, arg->child_shards, arg->pool_off, arg->initial_x, arg->initial_y, arg->running, arg->simplification, arg->layermaps, arg->layer_unmaps, arg->tiling_seg, arg->pass, arg->mingap, arg->minextent, arg->fraction, arg->prefilter, arg->postfilter, arg->filter, arg, arg->strategy, arg->compressed);
if (len < 0) { if (len < 0) {
int *err = &arg->err; int *err = &arg->err;
@@ -2908,7 +2968,7 @@ void *run_thread(void *vargs) {
return NULL; return NULL;
} }
int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpool, std::atomic<unsigned> *midx, std::atomic<unsigned> *midy, int &maxzoom, int minzoom, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, const char *tmpdir, double gamma, int full_detail, int low_detail, int min_detail, long long *meta_off, long long *pool_off, unsigned *initial_x, unsigned *initial_y, double simplification, double maxzoom_simplification, std::vector<std::map<std::string, layermap_entry>> &layermaps, const char *prefilter, const char *postfilter, std::map<std::string, attribute_op> const *attribute_accum, struct json_object *filter, std::vector<strategy> &strategies) { int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic<unsigned> *midx, std::atomic<unsigned> *midy, int &maxzoom, int minzoom, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, const char *tmpdir, double gamma, int full_detail, int low_detail, int min_detail, long long *pool_off, unsigned *initial_x, unsigned *initial_y, double simplification, double maxzoom_simplification, std::vector<std::map<std::string, layermap_entry>> &layermaps, const char *prefilter, const char *postfilter, std::map<std::string, attribute_op> const *attribute_accum, struct json_object *filter, std::vector<strategy> &strategies, int iz) {
last_progress = 0; last_progress = 0;
// The existing layermaps are one table per input thread. // The existing layermaps are one table per input thread.
@@ -2933,10 +2993,11 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
} }
int z; int z;
for (z = 0; z <= maxzoom; z++) { for (z = iz; z <= maxzoom; z++) {
std::atomic<long long> most(0); std::atomic<long long> most(0);
FILE *sub[TEMP_FILES]; compressor compressors[TEMP_FILES];
compressor *sub[TEMP_FILES];
int subfd[TEMP_FILES]; int subfd[TEMP_FILES];
for (size_t j = 0; j < TEMP_FILES; j++) { for (size_t j = 0; j < TEMP_FILES; j++) {
char geomname[strlen(tmpdir) + strlen("/geom.XXXXXXXX" XSTRINGIFY(INT_MAX)) + 1]; char geomname[strlen(tmpdir) + strlen("/geom.XXXXXXXX" XSTRINGIFY(INT_MAX)) + 1];
@@ -2947,11 +3008,13 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
perror(geomname); perror(geomname);
exit(EXIT_OPEN); exit(EXIT_OPEN);
} }
sub[j] = fopen_oflag(geomname, "wb", O_WRONLY | O_CLOEXEC); FILE *fp = fopen_oflag(geomname, "wb", O_WRONLY | O_CLOEXEC);
if (sub[j] == NULL) { if (fp == NULL) {
perror(geomname); perror(geomname);
exit(EXIT_OPEN); exit(EXIT_OPEN);
} }
compressors[j] = compressor(fp);
sub[j] = &compressors[j];
unlink(geomname); unlink(geomname);
} }
@@ -2970,7 +3033,7 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
} }
// XXX is it useful to divide further if we know we are skipping // XXX is it useful to divide further if we know we are skipping
// some zoom levels? Is it faster to have fewer CPUs working on // some zoom levels? Is it faster to have fewer CPUs working on
// sharding, but more deeply, or fewer CPUs, less deeply? // sharding, but more deeply, or more CPUs, less deeply?
if (threads > useful_threads) { if (threads > useful_threads) {
threads = useful_threads; threads = useful_threads;
} }
@@ -3055,7 +3118,6 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
atomic_strategy strategy; atomic_strategy strategy;
for (size_t thread = 0; thread < threads; thread++) { for (size_t thread = 0; thread < threads; thread++) {
args[thread].metabase = metabase;
args[thread].stringpool = stringpool; args[thread].stringpool = stringpool;
args[thread].min_detail = min_detail; args[thread].min_detail = min_detail;
args[thread].outdb = outdb; // locked with db_lock args[thread].outdb = outdb; // locked with db_lock
@@ -3092,7 +3154,6 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
args[thread].full_detail = full_detail; args[thread].full_detail = full_detail;
args[thread].low_detail = low_detail; args[thread].low_detail = low_detail;
args[thread].most = &most; // locked with var_lock args[thread].most = &most; // locked with var_lock
args[thread].meta_off = meta_off;
args[thread].pool_off = pool_off; args[thread].pool_off = pool_off;
args[thread].initial_x = initial_x; args[thread].initial_x = initial_x;
args[thread].initial_y = initial_y; args[thread].initial_y = initial_y;
@@ -3110,6 +3171,8 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
args[thread].wrote_zoom = -1; args[thread].wrote_zoom = -1;
args[thread].still_dropping = false; args[thread].still_dropping = false;
args[thread].strategy = &strategy; args[thread].strategy = &strategy;
args[thread].zoom = z;
args[thread].compressed = (z != iz);
if (pthread_create(&pthreads[thread], NULL, run_thread, &args[thread]) != 0) { if (pthread_create(&pthreads[thread], NULL, run_thread, &args[thread]) != 0) {
perror("pthread_create"); perror("pthread_create");
@@ -3188,7 +3251,7 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpo
exit(EXIT_CLOSE); exit(EXIT_CLOSE);
} }
} }
if (fclose(sub[j]) != 0) { if (sub[j]->fclose() != 0) {
perror("close subfile"); perror("close subfile");
exit(EXIT_CLOSE); exit(EXIT_CLOSE);
} }
+2 -2
View File
@@ -61,9 +61,9 @@ struct strategy {
strategy() = default; strategy() = default;
}; };
long long write_tile(char **geom, char *metabase, char *stringpool, unsigned *file_bbox, int z, unsigned x, unsigned y, int detail, int min_detail, int basezoom, sqlite3 *outdb, const char *outdir, double droprate, int buffer, const char *fname, FILE **geomfile, int file_minzoom, int file_maxzoom, double todo, char *geomstart, long long along, double gamma, int nlayers, std::atomic<strategy> *strategy); long long write_tile(char **geom, char *stringpool, unsigned *file_bbox, int z, unsigned x, unsigned y, int detail, int min_detail, int basezoom, sqlite3 *outdb, const char *outdir, double droprate, int buffer, const char *fname, FILE **geomfile, int file_minzoom, int file_maxzoom, double todo, char *geomstart, long long along, double gamma, int nlayers, std::atomic<strategy> *strategy);
int traverse_zooms(int *geomfd, off_t *geom_size, char *metabase, char *stringpool, std::atomic<unsigned> *midx, std::atomic<unsigned> *midy, int &maxzoom, int minzoom, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, const char *tmpdir, double gamma, int full_detail, int low_detail, int min_detail, long long *meta_off, long long *pool_off, unsigned *initial_x, unsigned *initial_y, double simplification, double maxzoom_simplification, std::vector<std::map<std::string, layermap_entry> > &layermap, const char *prefilter, const char *postfilter, std::map<std::string, attribute_op> const *attribute_accum, struct json_object *filter, std::vector<strategy> &strategies); int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic<unsigned> *midx, std::atomic<unsigned> *midy, int &maxzoom, int minzoom, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, const char *tmpdir, double gamma, int full_detail, int low_detail, int min_detail, long long *pool_off, unsigned *initial_x, unsigned *initial_y, double simplification, double maxzoom_simplification, std::vector<std::map<std::string, layermap_entry> > &layermap, const char *prefilter, const char *postfilter, std::map<std::string, attribute_op> const *attribute_accum, struct json_object *filter, std::vector<strategy> &strategies, int iz);
int manage_gap(unsigned long long index, unsigned long long *previndex, double scale, double gamma, double *gap); int manage_gap(unsigned long long index, unsigned long long *previndex, double scale, double gamma, double *gap);
+1 -1
View File
@@ -1,6 +1,6 @@
#ifndef VERSION_HPP #ifndef VERSION_HPP
#define VERSION_HPP #define VERSION_HPP
#define VERSION "v2.22.0" #define VERSION "v2.23.0"
#endif #endif