"tas" (type_and_string) => "sv" (serial_val)

This commit is contained in:
Erica Fischer
2024-01-23 10:21:45 -08:00
parent 0e2b2b8b7c
commit 4854c57e22
3 changed files with 53 additions and 53 deletions
+4 -4
View File
@@ -395,7 +395,7 @@ static void merge(struct mergelist *merges, size_t nmerges, unsigned char *map,
} }
struct sort_arg { struct sort_arg {
int task; int svk;
int cpus; int cpus;
long long indexpos; long long indexpos;
struct mergelist *merges; struct mergelist *merges;
@@ -404,8 +404,8 @@ struct sort_arg {
long long unit; long long unit;
int bytes; int bytes;
sort_arg(int task1, int cpus1, long long indexpos1, struct mergelist *merges1, int indexfd1, size_t nmerges1, long long unit1, int bytes1) sort_arg(int svk1, int cpus1, long long indexpos1, struct mergelist *merges1, int indexfd1, size_t nmerges1, long long unit1, int bytes1)
: task(task1), cpus(cpus1), indexpos(indexpos1), merges(merges1), indexfd(indexfd1), nmerges(nmerges1), unit(unit1), bytes(bytes1) { : svk(svk1), cpus(cpus1), indexpos(indexpos1), merges(merges1), indexfd(indexfd1), nmerges(nmerges1), unit(unit1), bytes(bytes1) {
} }
}; };
@@ -413,7 +413,7 @@ void *run_sort(void *v) {
struct sort_arg *a = (struct sort_arg *) v; struct sort_arg *a = (struct sort_arg *) v;
long long start; long long start;
for (start = a->task * a->unit; start < a->indexpos; start += a->unit * a->cpus) { for (start = a->svk * a->unit; start < a->indexpos; start += a->unit * a->cpus) {
long long end = start + a->unit; long long end = start + a->unit;
if (end > a->indexpos) { if (end > a->indexpos) {
end = a->indexpos; end = a->indexpos;
+21 -21
View File
@@ -206,11 +206,11 @@ void append_tile(std::string message, int z, unsigned x, unsigned y, std::map<st
} }
if (include.count(std::string(key)) || (!exclude_all && exclude.count(std::string(key)) == 0 && exclude_attributes.count(std::string(key)) == 0)) { if (include.count(std::string(key)) || (!exclude_all && exclude.count(std::string(key)) == 0 && exclude_attributes.count(std::string(key)) == 0)) {
serial_val tas; serial_val sv;
tas.type = type; sv.type = type;
tas.s = value; sv.s = value;
attributes.insert(std::pair<std::string, std::pair<mvt_value, serial_val>>(key, std::pair<mvt_value, serial_val>(val, tas))); attributes.insert(std::pair<std::string, std::pair<mvt_value, serial_val>>(key, std::pair<mvt_value, serial_val>(val, sv)));
key_order.push_back(key); key_order.push_back(key);
} }
@@ -253,14 +253,14 @@ void append_tile(std::string message, int z, unsigned x, unsigned y, std::map<st
attributes.erase(fa); attributes.erase(fa);
} }
serial_val tas; serial_val sv;
tas.type = outval.type; sv.type = outval.type;
tas.s = joinval; sv.s = joinval;
// Convert from double to int if the joined attribute is an integer // Convert from double to int if the joined attribute is an integer
outval = stringified_to_mvt_value(outval.type, joinval.c_str()); outval = stringified_to_mvt_value(outval.type, joinval.c_str());
attributes.insert(std::pair<std::string, std::pair<mvt_value, serial_val>>(joinkey, std::pair<mvt_value, serial_val>(outval, tas))); attributes.insert(std::pair<std::string, std::pair<mvt_value, serial_val>>(joinkey, std::pair<mvt_value, serial_val>(outval, sv)));
key_order.push_back(joinkey); key_order.push_back(joinkey);
} }
} }
@@ -821,7 +821,7 @@ void *join_worker(void *v) {
return NULL; return NULL;
} }
void dispatch_tasks(std::map<zxy, std::vector<std::string>> &tasks, std::vector<std::map<std::string, layermap_entry>> &layermaps, sqlite3 *outdb, const char *outdir, std::vector<std::string> &header, std::map<std::string, std::vector<std::string>> &mapping, std::set<std::string> &exclude, std::set<std::string> &include, int ifmatched, std::set<std::string> &keep_layers, std::set<std::string> &remove_layers, json_object *filter, struct tileset_reader *readers) { void dispatch_svks(std::map<zxy, std::vector<std::string>> &svks, std::vector<std::map<std::string, layermap_entry>> &layermaps, sqlite3 *outdb, const char *outdir, std::vector<std::string> &header, std::map<std::string, std::vector<std::string>> &mapping, std::set<std::string> &exclude, std::set<std::string> &include, int ifmatched, std::set<std::string> &keep_layers, std::set<std::string> &remove_layers, json_object *filter, struct tileset_reader *readers) {
pthread_t pthreads[CPUS]; pthread_t pthreads[CPUS];
std::vector<arg> args; std::vector<arg> args;
@@ -841,14 +841,14 @@ void dispatch_tasks(std::map<zxy, std::vector<std::string>> &tasks, std::vector<
} }
size_t count = 0; size_t count = 0;
// This isn't careful about distributing tasks evenly across CPUs, // This isn't careful about distributing svks evenly across CPUs,
// but, from testing, it actually takes a little longer to do // but, from testing, it actually takes a little longer to do
// the proper allocation than is saved by perfectly balanced threads. // the proper allocation than is saved by perfectly balanced threads.
for (auto ai = tasks.begin(); ai != tasks.end(); ++ai) { for (auto ai = svks.begin(); ai != svks.end(); ++ai) {
args[count].inputs.insert(*ai); args[count].inputs.insert(*ai);
count = (count + 1) % CPUS; count = (count + 1) % CPUS;
if (ai == tasks.begin()) { if (ai == svks.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); fflush(stderr);
@@ -975,7 +975,7 @@ void decode(struct tileset_reader *readers, std::map<std::string, layermap_entry
layermaps.push_back(std::map<std::string, layermap_entry>()); layermaps.push_back(std::map<std::string, layermap_entry>());
} }
std::map<zxy, std::vector<std::string>> tasks; std::map<zxy, std::vector<std::string>> svks;
double minlat = INT_MAX; double minlat = INT_MAX;
double minlon = INT_MAX; double minlon = INT_MAX;
double maxlat = INT_MIN; double maxlat = INT_MIN;
@@ -1014,14 +1014,14 @@ void decode(struct tileset_reader *readers, std::map<std::string, layermap_entry
if (current.first.z >= minzoom && current.first.z <= maxzoom) { if (current.first.z >= minzoom && current.first.z <= maxzoom) {
zxy tile = current.first; zxy tile = current.first;
if (tasks.count(tile) == 0) { if (svks.count(tile) == 0) {
tasks.insert(std::pair<zxy, std::vector<std::string>>(tile, std::vector<std::string>())); svks.insert(std::pair<zxy, std::vector<std::string>>(tile, std::vector<std::string>()));
} }
auto f = tasks.find(tile); auto f = svks.find(tile);
f->second.push_back(current.second); f->second.push_back(current.second);
} }
// Advance the tileset_reader that we just added as a task. // Advance the tileset_reader that we just added as a svk.
// The reason this prefetches is so the tileset_reader queue can be // The reason this prefetches is so the tileset_reader queue can be
// priority-ordered, so the one with the next relevant tile // priority-ordered, so the one with the next relevant tile
// is first in line. // is first in line.
@@ -1037,9 +1037,9 @@ void decode(struct tileset_reader *readers, std::map<std::string, layermap_entry
// Then this tile is done and we can safely run the output queue. // Then this tile is done and we can safely run the output queue.
if (readers == NULL || readers->zoom != current.first.z || readers->x != current.first.x || readers->y != current.first.y) { if (readers == NULL || readers->zoom != current.first.z || readers->x != current.first.x || readers->y != current.first.y) {
if (tasks.size() > 100 * CPUS) { if (svks.size() > 100 * CPUS) {
dispatch_tasks(tasks, layermaps, outdb, outdir, header, mapping, exclude, include, ifmatched, keep_layers, remove_layers, filter, readers); dispatch_svks(svks, layermaps, outdb, outdir, header, mapping, exclude, include, ifmatched, keep_layers, remove_layers, filter, readers);
tasks.clear(); svks.clear();
} }
} }
@@ -1067,7 +1067,7 @@ void decode(struct tileset_reader *readers, std::map<std::string, layermap_entry
st->minlat2 = min(minlat, st->minlat2); st->minlat2 = min(minlat, st->minlat2);
st->maxlat2 = max(maxlat, st->maxlat2); st->maxlat2 = max(maxlat, st->maxlat2);
dispatch_tasks(tasks, layermaps, outdb, outdir, header, mapping, exclude, include, ifmatched, keep_layers, remove_layers, filter, readers); dispatch_svks(svks, layermaps, outdb, outdir, header, mapping, exclude, include, ifmatched, keep_layers, remove_layers, filter, readers);
layermap = merge_layermaps(layermaps); layermap = merge_layermaps(layermaps);
struct tileset_reader *next; struct tileset_reader *next;
+28 -28
View File
@@ -546,8 +546,8 @@ struct partial {
struct partial_arg { struct partial_arg {
std::vector<struct partial> *partials = NULL; std::vector<struct partial> *partials = NULL;
int task = 0; int svk = 0;
int tasks = 0; int svks = 0;
drawvec *shared_nodes; drawvec *shared_nodes;
node *shared_nodes_map; node *shared_nodes_map;
@@ -672,7 +672,7 @@ void *partial_feature_worker(void *v) {
struct partial_arg *a = (struct partial_arg *) v; struct partial_arg *a = (struct partial_arg *) v;
std::vector<struct partial> *partials = a->partials; std::vector<struct partial> *partials = a->partials;
for (size_t i = a->task; i < (*partials).size(); i += a->tasks) { for (size_t i = a->svk; i < (*partials).size(); i += a->svks) {
double area = simplify_partial(&((*partials)[i]), *(a->shared_nodes), a->shared_nodes_map, a->nodepos); double area = simplify_partial(&((*partials)[i]), *(a->shared_nodes), a->shared_nodes_map, a->nodepos);
signed char t = (*partials)[i].t; signed char t = (*partials)[i].t;
@@ -1367,7 +1367,7 @@ long long choose_minextent(std::vector<long long> &extents, double f) {
} }
struct write_tile_args { struct write_tile_args {
struct task *tasks = NULL; struct svk *svks = NULL;
char *stringpool = NULL; char *stringpool = NULL;
int min_detail = 0; int min_detail = 0;
sqlite3 *outdb = NULL; sqlite3 *outdb = NULL;
@@ -2557,23 +2557,23 @@ long long write_tile(decompressor *geoms, std::atomic<long long> *geompos_in, ch
find_common_edges(partials, z, line_detail, simplification, maxzoom, merge_fraction); find_common_edges(partials, z, line_detail, simplification, maxzoom, merge_fraction);
} }
int tasks = ceil((double) CPUS / *running); int svks = ceil((double) CPUS / *running);
if (tasks < 1) { if (svks < 1) {
tasks = 1; svks = 1;
} }
pthread_t pthreads[tasks]; pthread_t pthreads[svks];
std::vector<partial_arg> args; std::vector<partial_arg> args;
args.resize(tasks); args.resize(svks);
for (int i = 0; i < tasks; i++) { for (int i = 0; i < svks; i++) {
args[i].task = i; args[i].svk = i;
args[i].tasks = tasks; args[i].svks = svks;
args[i].partials = &partials; args[i].partials = &partials;
args[i].shared_nodes = &shared_nodes; args[i].shared_nodes = &shared_nodes;
args[i].shared_nodes_map = shared_nodes_map; args[i].shared_nodes_map = shared_nodes_map;
args[i].nodepos = nodepos; args[i].nodepos = nodepos;
if (tasks > 1) { if (svks > 1) {
if (pthread_create(&pthreads[i], NULL, partial_feature_worker, &args[i]) != 0) { if (pthread_create(&pthreads[i], NULL, partial_feature_worker, &args[i]) != 0) {
perror("pthread_create"); perror("pthread_create");
exit(EXIT_PTHREAD); exit(EXIT_PTHREAD);
@@ -2583,8 +2583,8 @@ long long write_tile(decompressor *geoms, std::atomic<long long> *geompos_in, ch
} }
} }
if (tasks > 1) { if (svks > 1) {
for (int i = 0; i < tasks; i++) { for (int i = 0; i < svks; i++) {
void *retval; void *retval;
if (pthread_join(pthreads[i], &retval) != 0) { if (pthread_join(pthreads[i], &retval) != 0) {
@@ -3055,17 +3055,17 @@ long long write_tile(decompressor *geoms, std::atomic<long long> *geompos_in, ch
return -1; return -1;
} }
struct task { struct svk {
int fileno = 0; int fileno = 0;
struct task *next = NULL; struct svk *next = NULL;
}; };
void *run_thread(void *vargs) { void *run_thread(void *vargs) {
write_tile_args *arg = (write_tile_args *) vargs; write_tile_args *arg = (write_tile_args *) vargs;
struct task *task; struct svk *svk;
for (task = arg->tasks; task != NULL; task = task->next) { for (svk = arg->svks; svk != NULL; svk = svk->next) {
int j = task->fileno; int j = svk->fileno;
if (arg->geomfd[j] < 0) { if (arg->geomfd[j] < 0) {
// only one source file for zoom level 0 // only one source file for zoom level 0
@@ -3275,11 +3275,11 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic<
// Assign temporary files to threads // Assign temporary files to threads
std::vector<struct task> tasks; std::vector<struct svk> svks;
tasks.resize(TEMP_FILES); svks.resize(TEMP_FILES);
struct dispatch { struct dispatch {
struct task *tasks = NULL; struct svk *svks = NULL;
long long todo = 0; long long todo = 0;
struct dispatch *next = NULL; struct dispatch *next = NULL;
}; };
@@ -3288,7 +3288,7 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic<
struct dispatch *dispatch_head = &dispatches[0]; struct dispatch *dispatch_head = &dispatches[0];
for (size_t j = 0; j < threads; j++) { for (size_t j = 0; j < threads; j++) {
dispatches[j].tasks = NULL; dispatches[j].svks = NULL;
dispatches[j].todo = 0; dispatches[j].todo = 0;
if (j + 1 < threads) { if (j + 1 < threads) {
dispatches[j].next = &dispatches[j + 1]; dispatches[j].next = &dispatches[j + 1];
@@ -3302,9 +3302,9 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic<
continue; continue;
} }
tasks[j].fileno = j; svks[j].fileno = j;
tasks[j].next = dispatch_head->tasks; svks[j].next = dispatch_head->svks;
dispatch_head->tasks = &tasks[j]; dispatch_head->svks = &svks[j];
dispatch_head->todo += geom_size[j]; dispatch_head->todo += geom_size[j];
struct dispatch *here = dispatch_head; struct dispatch *here = dispatch_head;
@@ -3386,7 +3386,7 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic<
args[thread].attribute_accum = attribute_accum; args[thread].attribute_accum = attribute_accum;
args[thread].filter = filter; args[thread].filter = filter;
args[thread].tasks = dispatches[thread].tasks; args[thread].svks = dispatches[thread].svks;
args[thread].running = &running; args[thread].running = &running;
args[thread].pass = pass; args[thread].pass = pass;
args[thread].wrote_zoom = -1; args[thread].wrote_zoom = -1;