mirror of
https://github.com/simdjson/simdjson
synced 2026-06-08 17:27:07 +00:00
8e7d1a5f09
This creates a "document" class with only user-facing document state (no parser internals). - document: user-facing document state - document::iterator: iterator (equivalent of ParsedJsonIterator) - document::parser: parser state plus a "docked" document we parse into (equivalent of ParsedJson) Usage: ```c++ auto doc = simdjson::document::parse(buf, len); // less efficient but simplest ``` ```c++ simdjson::document::parser parser; // reusable parser parser.allocate_capacity(len); simdjson::document* doc = parser.parse(buf, len); // pointer to doc inside parser doc = parser.parse(buf2, len); // reuses all buffers and overwrites doc; more efficient ```
441 lines
16 KiB
C++
441 lines
16 KiB
C++
#ifndef SIMDJSON_JSONSTREAM_H
|
|
#define SIMDJSON_JSONSTREAM_H
|
|
|
|
#include <algorithm>
|
|
#include <limits>
|
|
#include <stdexcept>
|
|
#include <thread>
|
|
#include "simdjson/isadetection.h"
|
|
#include "simdjson/padded_string.h"
|
|
#include "simdjson/simdjson.h"
|
|
#include "simdjson/stage1_find_marks.h"
|
|
#include "simdjson/stage2_build_tape.h"
|
|
#include "jsoncharutils.h"
|
|
|
|
|
|
namespace simdjson {
|
|
/*************************************************************************************
|
|
* The main motivation for this piece of software is to achieve maximum speed
|
|
*and offer
|
|
* good quality of life while parsing files containing multiple JSON documents.
|
|
*
|
|
* Since we want to offer flexibility and not restrict ourselves to a specific
|
|
*file
|
|
* format, we support any file that contains any valid JSON documents separated
|
|
*by one
|
|
* or more character that is considered a whitespace by the JSON spec.
|
|
* Namely: space, nothing, linefeed, carriage return, horizontal tab.
|
|
* Anything that is not whitespace will be parsed as a JSON document and could
|
|
*lead
|
|
* to failure.
|
|
*
|
|
* To offer maximum parsing speed, our implementation processes the data inside
|
|
*the
|
|
* buffer by batches and their size is defined by the parameter "batch_size".
|
|
* By loading data in batches, we can optimize the time spent allocating data in
|
|
*the
|
|
* parser and can also open the possibility of multi-threading.
|
|
* The batch_size must be at least as large as the biggest document in the file,
|
|
*but
|
|
* not too large in order to submerge the chached memory. We found that 1MB is
|
|
* somewhat a sweet spot for now. Eventually, this batch_size could be fully
|
|
* automated and be optimal at all times.
|
|
************************************************************************************/
|
|
/**
|
|
* The template parameter (string_container) must
|
|
* support the data() and size() methods, returning a pointer
|
|
* to a char* and to the number of bytes respectively.
|
|
* The simdjson parser may read up to SIMDJSON_PADDING bytes beyond the end
|
|
* of the string, so if you do not use a padded_string container,
|
|
* you have the responsability to overallocated. If you fail to
|
|
* do so, your software may crash if you cross a page boundary,
|
|
* and you should expect memory checkers to object.
|
|
* Most users should use a simdjson::padded_string.
|
|
*/
|
|
template <class string_container = padded_string> class JsonStream {
|
|
public:
|
|
/* Create a JsonStream object that can be used to parse sequentially the valid
|
|
* JSON documents found in the buffer "buf".
|
|
*
|
|
* The batch_size must be at least as large as the biggest document in the
|
|
* file, but
|
|
* not too large to submerge the cached memory. We found that 1MB is
|
|
* somewhat a sweet spot for now.
|
|
*
|
|
* The user is expected to call the following json_parse method to parse the
|
|
* next
|
|
* valid JSON document found in the buffer. This method can and is expected
|
|
* to be
|
|
* called in a loop.
|
|
*
|
|
* Various methods are offered to keep track of the status, like
|
|
* get_current_buffer_loc,
|
|
* get_n_parsed_docs, get_n_bytes_parsed, etc.
|
|
*
|
|
* */
|
|
JsonStream(const string_container &s, size_t batch_size = 1000000);
|
|
|
|
~JsonStream();
|
|
|
|
/* Parse the next document found in the buffer previously given to JsonStream.
|
|
|
|
* The content should be a valid JSON document encoded as UTF-8. If there is a
|
|
* UTF-8 BOM, the caller is responsible for omitting it, UTF-8 BOM are
|
|
* discouraged.
|
|
*
|
|
* You do NOT need to pre-allocate a parser. This function takes care of
|
|
* pre-allocating a capacity defined by the batch_size defined when creating
|
|
the
|
|
* JsonStream object.
|
|
*
|
|
* The function returns simdjson::SUCCESS_AND_HAS_MORE (an integer = 1) in
|
|
case
|
|
* of success and indicates that the buffer still contains more data to be
|
|
parsed,
|
|
* meaning this function can be called again to return the next JSON document
|
|
* after this one.
|
|
*
|
|
* The function returns simdjson::SUCCESS (as integer = 0) in case of success
|
|
* and indicates that the buffer has successfully been parsed to the end.
|
|
* Every document it contained has been parsed without error.
|
|
*
|
|
* The function returns an error code from simdjson/simdjson.h in case of
|
|
failure
|
|
* such as simdjson::CAPACITY, simdjson::MEMALLOC, simdjson::DEPTH_ERROR and
|
|
so forth;
|
|
* the simdjson::error_message function converts these error codes into a
|
|
* string).
|
|
*
|
|
* You can also check validity by calling parser.is_valid(). The same parser
|
|
can
|
|
* and should be reused for the other documents in the buffer. */
|
|
int json_parse(document::parser &parser);
|
|
|
|
/* Returns the location (index) of where the next document should be in the
|
|
* buffer.
|
|
* Can be used for debugging, it tells the user the position of the end of the
|
|
* last
|
|
* valid JSON document parsed*/
|
|
inline size_t get_current_buffer_loc() const { return current_buffer_loc; }
|
|
|
|
/* Returns the total amount of complete documents parsed by the JsonStream,
|
|
* in the current buffer, at the given time.*/
|
|
inline size_t get_n_parsed_docs() const { return n_parsed_docs; }
|
|
|
|
/* Returns the total amount of data (in bytes) parsed by the JsonStream,
|
|
* in the current buffer, at the given time.*/
|
|
inline size_t get_n_bytes_parsed() const { return n_bytes_parsed; }
|
|
|
|
private:
|
|
inline const char *buf() const { return str.data() + str_start; }
|
|
|
|
inline void advance(size_t offset) { str_start += offset; }
|
|
|
|
inline size_t remaining() const { return str.size() - str_start; }
|
|
|
|
const string_container &str;
|
|
size_t _batch_size; // this is actually variable!
|
|
size_t str_start{0};
|
|
size_t next_json{0};
|
|
bool load_next_batch{true};
|
|
size_t current_buffer_loc{0};
|
|
#ifdef SIMDJSON_THREADS_ENABLED
|
|
size_t last_json_buffer_loc{0};
|
|
#endif
|
|
size_t n_parsed_docs{0};
|
|
size_t n_bytes_parsed{0};
|
|
#ifdef SIMDJSON_THREADS_ENABLED
|
|
int stage1_is_ok_thread{0};
|
|
std::thread stage_1_thread;
|
|
document::parser parser_thread;
|
|
#endif
|
|
}; // end of class JsonStream
|
|
|
|
/* This algorithm is used to quickly identify the buffer position of
|
|
* the last JSON document inside the current batch.
|
|
*
|
|
* It does its work by finding the last pair of structural characters
|
|
* that represent the end followed by the start of a document.
|
|
*
|
|
* Simply put, we iterate over the structural characters, starting from
|
|
* the end. We consider that we found the end of a JSON document when the
|
|
* first element of the pair is NOT one of these characters: '{' '[' ';' ','
|
|
* and when the second element is NOT one of these characters: '}' '}' ';' ','.
|
|
*
|
|
* This simple comparison works most of the time, but it does not cover cases
|
|
* where the batch's structural indexes contain a perfect amount of documents.
|
|
* In such a case, we do not have access to the structural index which follows
|
|
* the last document, therefore, we do not have access to the second element in
|
|
* the pair, and means that we cannot identify the last document. To fix this
|
|
* issue, we keep a count of the open and closed curly/square braces we found
|
|
* while searching for the pair. When we find a pair AND the count of open and
|
|
* closed curly/square braces is the same, we know that we just passed a
|
|
* complete
|
|
* document, therefore the last json buffer location is the end of the batch
|
|
* */
|
|
inline size_t find_last_json_buf_idx(const char *buf, size_t size,
|
|
const document::parser &parser) {
|
|
// this function can be generally useful
|
|
if (parser.n_structural_indexes == 0)
|
|
return 0;
|
|
auto last_i = parser.n_structural_indexes - 1;
|
|
if (parser.structural_indexes[last_i] == size) {
|
|
if (last_i == 0)
|
|
return 0;
|
|
last_i = parser.n_structural_indexes - 2;
|
|
}
|
|
auto arr_cnt = 0;
|
|
auto obj_cnt = 0;
|
|
for (auto i = last_i; i > 0; i--) {
|
|
auto idxb = parser.structural_indexes[i];
|
|
switch (buf[idxb]) {
|
|
case ':':
|
|
case ',':
|
|
continue;
|
|
case '}':
|
|
obj_cnt--;
|
|
continue;
|
|
case ']':
|
|
arr_cnt--;
|
|
continue;
|
|
case '{':
|
|
obj_cnt++;
|
|
break;
|
|
case '[':
|
|
arr_cnt++;
|
|
break;
|
|
}
|
|
auto idxa = parser.structural_indexes[i - 1];
|
|
switch (buf[idxa]) {
|
|
case '{':
|
|
case '[':
|
|
case ':':
|
|
case ',':
|
|
continue;
|
|
}
|
|
if (!arr_cnt && !obj_cnt) {
|
|
return last_i + 1;
|
|
}
|
|
return i;
|
|
}
|
|
return 0;
|
|
}
|
|
|
|
// Everything in the following anonymous namespace should go.
|
|
// It is a hack.
|
|
namespace {
|
|
|
|
typedef int (*stage1_functype)(const char *buf, size_t len,
|
|
document::parser &parser, bool streaming);
|
|
typedef int (*stage2_functype)(const char *buf, size_t len,
|
|
document::parser &parser, size_t &next_json);
|
|
|
|
stage1_functype best_stage1;
|
|
stage2_functype best_stage2;
|
|
|
|
//// TODO: generalize this set of functions. We don't want to have a copy in
|
|
/// jsonparser.cpp
|
|
void find_the_best_supported_implementation() {
|
|
uint32_t supports = simdjson::detect_supported_architectures();
|
|
// Order from best to worst (within architecture)
|
|
#ifdef IS_X86_64
|
|
constexpr uint32_t haswell_flags =
|
|
simdjson::instruction_set::AVX2 | simdjson::instruction_set::PCLMULQDQ |
|
|
simdjson::instruction_set::BMI1 | simdjson::instruction_set::BMI2;
|
|
constexpr uint32_t westmere_flags =
|
|
simdjson::instruction_set::SSE42 | simdjson::instruction_set::PCLMULQDQ;
|
|
if ((haswell_flags & supports) == haswell_flags) {
|
|
best_stage1 =
|
|
simdjson::find_structural_bits<simdjson::Architecture::HASWELL>;
|
|
best_stage2 = simdjson::unified_machine<simdjson::Architecture::HASWELL>;
|
|
return;
|
|
}
|
|
if ((westmere_flags & supports) == westmere_flags) {
|
|
best_stage1 =
|
|
simdjson::find_structural_bits<simdjson::Architecture::WESTMERE>;
|
|
best_stage2 = simdjson::unified_machine<simdjson::Architecture::WESTMERE>;
|
|
return;
|
|
}
|
|
#endif
|
|
#ifdef IS_ARM64
|
|
if (supports & instruction_set::NEON) {
|
|
best_stage1 = simdjson::find_structural_bits<Architecture::ARM64>;
|
|
best_stage2 = simdjson::unified_machine<Architecture::ARM64>;
|
|
return;
|
|
}
|
|
#endif
|
|
// we throw an exception since this should not be recoverable
|
|
throw new std::runtime_error("unsupported architecture");
|
|
}
|
|
} // anonymous namespace
|
|
|
|
template <class string_container>
|
|
JsonStream<string_container>::JsonStream(const string_container &s,
|
|
size_t batchSize)
|
|
: str(s), _batch_size(batchSize) {
|
|
find_the_best_supported_implementation();
|
|
}
|
|
|
|
template <class string_container> JsonStream<string_container>::~JsonStream() {
|
|
#ifdef SIMDJSON_THREADS_ENABLED
|
|
if (stage_1_thread.joinable()) {
|
|
stage_1_thread.join();
|
|
}
|
|
#endif
|
|
}
|
|
|
|
#ifdef SIMDJSON_THREADS_ENABLED
|
|
|
|
// threaded version of json_parse
|
|
// todo: simplify this code further
|
|
template <class string_container>
|
|
int JsonStream<string_container>::json_parse(document::parser &parser) {
|
|
if (unlikely(parser.capacity() == 0)) {
|
|
const bool allocok = parser.allocate_capacity(_batch_size);
|
|
if (!allocok) {
|
|
parser.error_code = simdjson::MEMALLOC;
|
|
return parser.error_code;
|
|
}
|
|
} else if (unlikely(parser.capacity() < _batch_size)) {
|
|
parser.error_code = simdjson::CAPACITY;
|
|
return parser.error_code;
|
|
}
|
|
if (unlikely(parser_thread.capacity() < _batch_size)) {
|
|
const bool allocok_thread = parser_thread.allocate_capacity(_batch_size);
|
|
if (!allocok_thread) {
|
|
parser.error_code = simdjson::MEMALLOC;
|
|
return parser.error_code;
|
|
}
|
|
}
|
|
if (unlikely(load_next_batch)) {
|
|
// First time loading
|
|
if (!stage_1_thread.joinable()) {
|
|
_batch_size = (std::min)(_batch_size, remaining());
|
|
_batch_size = trimmed_length_safe_utf8((const char *)buf(), _batch_size);
|
|
if (_batch_size == 0) {
|
|
parser.error_code = simdjson::UTF8_ERROR;
|
|
return parser.error_code;
|
|
}
|
|
int stage1_is_ok = best_stage1(buf(), _batch_size, parser, true);
|
|
if (stage1_is_ok != simdjson::SUCCESS) {
|
|
parser.error_code = stage1_is_ok;
|
|
return parser.error_code;
|
|
}
|
|
size_t last_index = find_last_json_buf_idx(buf(), _batch_size, parser);
|
|
if (last_index == 0) {
|
|
if (parser.n_structural_indexes == 0) {
|
|
parser.error_code = simdjson::EMPTY;
|
|
return parser.error_code;
|
|
}
|
|
} else {
|
|
parser.n_structural_indexes = last_index + 1;
|
|
}
|
|
}
|
|
// the second thread is running or done.
|
|
else {
|
|
stage_1_thread.join();
|
|
if (stage1_is_ok_thread != simdjson::SUCCESS) {
|
|
parser.error_code = stage1_is_ok_thread;
|
|
return parser.error_code;
|
|
}
|
|
std::swap(parser.structural_indexes, parser_thread.structural_indexes);
|
|
parser.n_structural_indexes = parser_thread.n_structural_indexes;
|
|
advance(last_json_buffer_loc);
|
|
n_bytes_parsed += last_json_buffer_loc;
|
|
}
|
|
// let us decide whether we will start a new thread
|
|
if (remaining() - _batch_size > 0) {
|
|
last_json_buffer_loc =
|
|
parser.structural_indexes[find_last_json_buf_idx(buf(), _batch_size, parser)];
|
|
_batch_size = (std::min)(_batch_size, remaining() - last_json_buffer_loc);
|
|
if (_batch_size > 0) {
|
|
_batch_size = trimmed_length_safe_utf8(
|
|
(const char *)(buf() + last_json_buffer_loc), _batch_size);
|
|
if (_batch_size == 0) {
|
|
parser.error_code = simdjson::UTF8_ERROR;
|
|
return parser.error_code;
|
|
}
|
|
// let us capture read-only variables
|
|
const char *const b = buf() + last_json_buffer_loc;
|
|
const size_t bs = _batch_size;
|
|
// we call the thread on a lambda that will update
|
|
// this->stage1_is_ok_thread
|
|
// there is only one thread that may write to this value
|
|
stage_1_thread = std::thread([this, b, bs] {
|
|
this->stage1_is_ok_thread = best_stage1(b, bs, this->parser_thread, true);
|
|
});
|
|
}
|
|
}
|
|
next_json = 0;
|
|
load_next_batch = false;
|
|
} // load_next_batch
|
|
int res = best_stage2(buf(), remaining(), parser, next_json);
|
|
if (res == simdjson::SUCCESS_AND_HAS_MORE) {
|
|
n_parsed_docs++;
|
|
current_buffer_loc = parser.structural_indexes[next_json];
|
|
load_next_batch = (current_buffer_loc == last_json_buffer_loc);
|
|
} else if (res == simdjson::SUCCESS) {
|
|
n_parsed_docs++;
|
|
if (remaining() > _batch_size) {
|
|
current_buffer_loc = parser.structural_indexes[next_json - 1];
|
|
load_next_batch = true;
|
|
res = simdjson::SUCCESS_AND_HAS_MORE;
|
|
}
|
|
}
|
|
return res;
|
|
}
|
|
|
|
#else // SIMDJSON_THREADS_ENABLED
|
|
|
|
// single-threaded version of json_parse
|
|
template <class string_container>
|
|
int JsonStream<string_container>::json_parse(document::parser &parser) {
|
|
if (unlikely(parser.capacity() == 0)) {
|
|
const bool allocok = parser.allocate_capacity(_batch_size);
|
|
if (!allocok) {
|
|
return parser.on_error(MEMALLOC);
|
|
}
|
|
} else if (unlikely(parser.capacity() < _batch_size)) {
|
|
return parser.on_error(CAPACITY);
|
|
}
|
|
if (unlikely(load_next_batch)) {
|
|
advance(current_buffer_loc);
|
|
n_bytes_parsed += current_buffer_loc;
|
|
_batch_size = (std::min)(_batch_size, remaining());
|
|
_batch_size = trimmed_length_safe_utf8((const char *)buf(), _batch_size);
|
|
auto stage1_is_ok = (ErrorValues)best_stage1(buf(), _batch_size, parser, true);
|
|
if (stage1_is_ok != simdjson::SUCCESS) {
|
|
return parser.on_error(stage1_is_ok);
|
|
}
|
|
size_t last_index = find_last_json_buf_idx(buf(), _batch_size, parser);
|
|
if (last_index == 0) {
|
|
if (parser.n_structural_indexes == 0) {
|
|
return parser.on_error(EMPTY);
|
|
}
|
|
} else {
|
|
parser.n_structural_indexes = last_index + 1;
|
|
}
|
|
load_next_batch = false;
|
|
} // load_next_batch
|
|
int res = best_stage2(buf(), remaining(), parser, next_json);
|
|
if (likely(res == simdjson::SUCCESS_AND_HAS_MORE)) {
|
|
n_parsed_docs++;
|
|
current_buffer_loc = parser.structural_indexes[next_json];
|
|
} else if (res == simdjson::SUCCESS) {
|
|
n_parsed_docs++;
|
|
if (remaining() > _batch_size) {
|
|
current_buffer_loc = parser.structural_indexes[next_json - 1];
|
|
next_json = 1;
|
|
load_next_batch = true;
|
|
res = simdjson::SUCCESS_AND_HAS_MORE;
|
|
}
|
|
} else {
|
|
printf("E\n");
|
|
}
|
|
return res;
|
|
}
|
|
#endif // SIMDJSON_THREADS_ENABLED
|
|
|
|
} // end of namespace simdjson
|
|
#endif // SIMDJSON_JSONSTREAM_H
|