FlatGeobuf multithreading [#2]

This commit is contained in:
Brandon Liu
2022-03-27 19:14:27 +08:00
parent 245570feff
commit 6e6cd29399
+124 -33
View File
@@ -5,6 +5,7 @@
#include "flatgeobuf/feature_generated.h" #include "flatgeobuf/feature_generated.h"
#include "flatgeobuf/header_generated.h" #include "flatgeobuf/header_generated.h"
#include "milo/dtoa_milo.h" #include "milo/dtoa_milo.h"
#include "main.hpp"
using namespace std; using namespace std;
@@ -93,22 +94,8 @@ drawvec readGeometry(const FlatGeobuf::Geometry *geometry, FlatGeobuf::GeometryT
} }
} }
void parse_flatgeobuf(std::vector<struct serialization_state> *sst, const char *src, size_t len, int layer, std::string layername) { void readFeature(const FlatGeobuf::Feature *feature, FlatGeobuf::GeometryType h_geometry_type, const vector<string> &h_column_names, const vector<FlatGeobuf::ColumnType> &h_column_types, struct serialization_state *sst, int layer, std::string layername) {
auto header_size = flatbuffers::GetPrefixedSize((const uint8_t *)src + 8); drawvec dv = readGeometry(feature->geometry(), h_geometry_type);
auto header = FlatGeobuf::GetSizePrefixedHeader(src + 8);
auto features_count = header->features_count();
auto node_size = header->index_node_size();
vector<string> h_column_names;
vector<FlatGeobuf::ColumnType> h_column_types;
for (size_t i = 0; i < header->columns()->size(); i++) {
h_column_names.push_back(header->columns()->Get(i)->name()->c_str());
h_column_types.push_back(header->columns()->Get(i)->type());
}
auto h_geometry_type = header->geometry_type();
int drawvec_type = -1; int drawvec_type = -1;
@@ -132,30 +119,16 @@ void parse_flatgeobuf(std::vector<struct serialization_state> *sst, const char *
exit(EXIT_FAILURE); exit(EXIT_FAILURE);
} }
int index_size = 0;
if (node_size > 0) {
index_size = PackedRTreeSize(features_count,node_size);
}
const char* start = src + 8 + 4 + header_size + index_size;
while (start < src + len) {
serial_feature sf; serial_feature sf;
auto my_sst = &(*sst)[0];
auto feature_size = flatbuffers::GetPrefixedSize((const uint8_t *)start);
auto feature = FlatGeobuf::GetSizePrefixedFeature(start);
drawvec dv = readGeometry(feature->geometry(), h_geometry_type);
sf.layer = layer; sf.layer = layer;
sf.layername = layername; sf.layername = layername;
sf.segment = my_sst->segment; sf.segment = sst->segment;
sf.has_id = false; sf.has_id = false;
sf.has_tippecanoe_minzoom = false; sf.has_tippecanoe_minzoom = false;
sf.has_tippecanoe_maxzoom = false; sf.has_tippecanoe_maxzoom = false;
sf.feature_minzoom = false; sf.feature_minzoom = false;
sf.seq = (*my_sst->layer_seq); sf.seq = (*sst->layer_seq);
sf.geometry = dv; sf.geometry = dv;
sf.t = drawvec_type; sf.t = drawvec_type;
@@ -207,7 +180,125 @@ void parse_flatgeobuf(std::vector<struct serialization_state> *sst, const char *
sf.full_keys = full_keys; sf.full_keys = full_keys;
sf.full_values = full_values; sf.full_values = full_values;
serialize_feature(my_sst, sf); serialize_feature(sst, sf);
}
struct queued_feature {
const FlatGeobuf::Feature *feature = NULL;
FlatGeobuf::GeometryType h_geometry_type = FlatGeobuf::GeometryType::Unknown;
const vector<std::string> *h_column_names = NULL;
const vector<FlatGeobuf::ColumnType> *h_column_types = NULL;
vector<struct serialization_state> *sst = NULL;
int layer = 0;
std::string layername = "";
};
static std::vector<queued_feature> feature_queue;
struct queue_run_arg {
size_t start;
size_t end;
size_t segment;
queue_run_arg(size_t start1, size_t end1, size_t segment1)
: start(start1), end(end1), segment(segment1) {
}
};
void *fgb_run_parse_feature(void *v) {
struct queue_run_arg *qra = (struct queue_run_arg *) v;
for (size_t i = qra->start; i < qra->end; i++) {
struct queued_feature &qf = feature_queue[i];
readFeature(qf.feature, qf.h_geometry_type, *qf.h_column_names, *qf.h_column_types, &(*qf.sst)[qra->segment], qf.layer, qf.layername);
}
return NULL;
}
void fgbRunQueue() {
if (feature_queue.size() == 0) {
return;
}
std::vector<struct queue_run_arg> qra;
std::vector<pthread_t> pthreads;
pthreads.resize(CPUS);
for (size_t i = 0; i < CPUS; i++) {
*((*(feature_queue[0].sst))[i].layer_seq) = *((*(feature_queue[0].sst))[0].layer_seq) + feature_queue.size() * i / CPUS;
qra.push_back(queue_run_arg(
feature_queue.size() * i / CPUS,
feature_queue.size() * (i + 1) / CPUS,
i));
}
for (size_t i = 0; i < CPUS; i++) {
if (pthread_create(&pthreads[i], NULL, fgb_run_parse_feature, &qra[i]) != 0) {
perror("pthread_create");
exit(EXIT_FAILURE);
}
}
for (size_t i = 0; i < CPUS; i++) {
void *retval;
if (pthread_join(pthreads[i], &retval) != 0) {
perror("pthread_join");
}
}
// Lack of atomicity is OK, since we are single-threaded again here
long long was = *((*(feature_queue[0].sst))[CPUS - 1].layer_seq);
*((*(feature_queue[0].sst))[0].layer_seq) = was;
feature_queue.clear();
}
void queueFeature(const FlatGeobuf::Feature *feature, FlatGeobuf::GeometryType h_geometry_type, const vector<string> &h_column_names, const vector<FlatGeobuf::ColumnType> &h_column_types, std::vector<struct serialization_state> *sst, int layer, string layername) {
struct queued_feature qf;
qf.feature = feature;
qf.h_geometry_type = h_geometry_type;
qf.h_column_names = &h_column_names;
qf.h_column_types = &h_column_types;
qf.sst = sst;
qf.layer = layer;
qf.layername = layername;
feature_queue.push_back(qf);
if (feature_queue.size() > CPUS * 500) {
fgbRunQueue();
}
}
void parse_flatgeobuf(std::vector<struct serialization_state> *sst, const char *src, size_t len, int layer, std::string layername) {
auto header_size = flatbuffers::GetPrefixedSize((const uint8_t *)src + 8);
auto header = FlatGeobuf::GetSizePrefixedHeader(src + 8);
auto features_count = header->features_count();
auto node_size = header->index_node_size();
vector<string> h_column_names;
vector<FlatGeobuf::ColumnType> h_column_types;
for (size_t i = 0; i < header->columns()->size(); i++) {
h_column_names.push_back(header->columns()->Get(i)->name()->c_str());
h_column_types.push_back(header->columns()->Get(i)->type());
}
auto h_geometry_type = header->geometry_type();
int index_size = 0;
if (node_size > 0) {
index_size = PackedRTreeSize(features_count,node_size);
}
const char* start = src + 8 + 4 + header_size + index_size;
while (start < src + len) {
auto feature_size = flatbuffers::GetPrefixedSize((const uint8_t *)start);
auto feature = FlatGeobuf::GetSizePrefixedFeature(start);
queueFeature(feature, h_geometry_type, h_column_names, h_column_types, sst, layer, layername);
start += 4 + feature_size; start += 4 + feature_size;
} }