Coverage Report

Created: 2026-09-30 18:42

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/json/new_json_reader.h
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#pragma once
19
20
#include <rapidjson/allocators.h>
21
#include <rapidjson/document.h>
22
#include <rapidjson/encodings.h>
23
#include <rapidjson/rapidjson.h>
24
#include <simdjson/common_defs.h>
25
#include <simdjson/simdjson.h> // IWYU pragma: keep
26
27
#include <memory>
28
#include <string>
29
#include <string_view>
30
#include <unordered_map>
31
#include <unordered_set>
32
#include <vector>
33
34
#include "common/status.h"
35
#include "core/custom_allocator.h"
36
#include "core/string_ref.h"
37
#include "core/types.h"
38
#include "exprs/json_functions.h"
39
#include "format/line_reader.h"
40
#include "format/table/table_format_reader.h"
41
#include "io/file_factory.h"
42
#include "io/fs/file_reader_writer_fwd.h"
43
#include "runtime/runtime_profile.h"
44
#include "util/decompressor.h"
45
46
namespace simdjson::fallback::ondemand {
47
class object;
48
} // namespace simdjson::fallback::ondemand
49
50
namespace doris {
51
class SlotDescriptor;
52
class RuntimeState;
53
class TFileRangeDesc;
54
class TFileScanRangeParams;
55
56
namespace io {
57
class FileSystem;
58
struct IOContext;
59
} // namespace io
60
61
struct ScannerCounter;
62
class Block;
63
class IColumn;
64
class NewPlainTextLineReader;
65
66
namespace json_reader_detail {
67
Status append_null_for_malformed_json(Block& block);
68
void truncate_block_to_rows(Block& block, size_t num_rows);
69
void pop_back_last_inserted_value(Block& block, size_t column_index);
70
} // namespace json_reader_detail
71
72
/// JSON-specific initialization context.
73
/// Extends ReaderInitContext with default value context (unique to JSON reader).
74
struct JsonInitContext final : public ReaderInitContext {
75
    const std::unordered_map<std::string, VExprContextSPtr>* col_default_value_ctx = nullptr;
76
    bool is_load = false;
77
};
78
79
class NewJsonReader : public TableFormatReader {
80
    ENABLE_FACTORY_CREATOR(NewJsonReader);
81
82
public:
83
    NewJsonReader(RuntimeState* state, RuntimeProfile* profile, ScannerCounter* counter,
84
                  const TFileScanRangeParams& params, const TFileRangeDesc& range,
85
                  const std::vector<SlotDescriptor*>& file_slot_descs, bool* scanner_eof,
86
                  size_t batch_size, io::IOContext* io_ctx,
87
                  std::shared_ptr<io::IOContext> io_ctx_holder = nullptr);
88
89
    NewJsonReader(RuntimeProfile* profile, const TFileScanRangeParams& params,
90
                  const TFileRangeDesc& range, const std::vector<SlotDescriptor*>& file_slot_descs,
91
                  size_t batch_size, io::IOContext* io_ctx,
92
                  std::shared_ptr<io::IOContext> io_ctx_holder = nullptr);
93
    ~NewJsonReader() override;
94
95
    Status init_reader(
96
            const std::unordered_map<std::string, VExprContextSPtr>& col_default_value_ctx,
97
            bool is_load);
98
99
    Status _do_get_next_block(Block* block, size_t* read_rows, bool* eof) override;
100
    Status _get_columns_impl(std::unordered_map<std::string, DataTypePtr>* name_to_type) override;
101
102
    // Row-based readers control throughput via row count, not byte budget.
103
    // The FileScanner's AdaptiveBlockSizePredictor converts the byte budget
104
    // into a predicted row count and calls set_batch_size() with it.
105
    void set_batch_size(size_t batch_size) override;
106
6
    size_t get_batch_size() const override { return _batch_size; }
107
108
    Status init_schema_reader() override;
109
    Status get_parsed_schema(std::vector<std::string>* col_names,
110
                             std::vector<DataTypePtr>* col_types) override;
111
112
protected:
113
    // ---- Unified init_reader(ReaderInitContext*) overrides ----
114
    Status _open_file_reader(ReaderInitContext* ctx) override;
115
    Status _do_init_reader(ReaderInitContext* ctx) override;
116
117
    void _collect_profile_before_close() override;
118
119
private:
120
    Status _get_range_params();
121
    void _init_system_properties();
122
    void _init_file_description();
123
    Status _open_file_reader(bool need_schema);
124
    Status _open_line_reader();
125
    Status _parse_jsonpath_and_json_root();
126
127
    Status _read_json_column(RuntimeState* state, Block& block,
128
                             const std::vector<SlotDescriptor*>& slot_descs, bool* is_empty_row,
129
                             bool* eof);
130
131
    Status _read_one_message(DorisUniqueBufferPtr<uint8_t>* file_buf, size_t* read_size);
132
133
    // StreamLoadPipe::read_one_message only reads a portion of the data when stream loading with a chunked transfer HTTP request.
134
    // Need to read all the data before performing JSON parsing.
135
    Status _read_one_message_from_pipe(DorisUniqueBufferPtr<uint8_t>* file_buf, size_t* read_size);
136
137
    // simdjson, replace none simdjson function if it is ready
138
    Status _simdjson_init_reader();
139
    Status _simdjson_parse_json(size_t* size, bool* is_empty_row, bool* eof,
140
                                simdjson::error_code* error);
141
    Status _get_json_value(size_t* size, bool* eof, simdjson::error_code* error,
142
                           bool* is_empty_row);
143
    Status _judge_empty_row(size_t size, bool eof, bool* is_empty_row);
144
145
    Status _handle_simdjson_error(simdjson::simdjson_error& error, Block& block, size_t num_rows,
146
                                  bool* eof);
147
148
    Status _simdjson_handle_simple_json(RuntimeState* state, Block& block,
149
                                        const std::vector<SlotDescriptor*>& slot_descs,
150
                                        bool* is_empty_row, bool* eof);
151
152
    Status _simdjson_handle_simple_json_write_columns(
153
            Block& block, const std::vector<SlotDescriptor*>& slot_descs, bool* is_empty_row,
154
            bool* eof);
155
156
    Status _simdjson_handle_flat_array_complex_json(RuntimeState* state, Block& block,
157
                                                    const std::vector<SlotDescriptor*>& slot_descs,
158
                                                    bool* is_empty_row, bool* eof);
159
160
    Status _simdjson_handle_flat_array_complex_json_write_columns(
161
            Block& block, const std::vector<SlotDescriptor*>& slot_descs, bool* is_empty_row,
162
            bool* eof);
163
164
    Status _simdjson_handle_nested_complex_json(RuntimeState* state, Block& block,
165
                                                const std::vector<SlotDescriptor*>& slot_descs,
166
                                                bool* is_empty_row, bool* eof);
167
168
    Status _simdjson_set_column_value(simdjson::ondemand::object* value, Block& block,
169
                                      const std::vector<SlotDescriptor*>& slot_descs, bool* valid);
170
171
    template <bool use_string_cache>
172
    Status _simdjson_write_data_to_column(simdjson::ondemand::value& value,
173
                                          const DataTypePtr& type_desc, IColumn* column_ptr,
174
                                          const std::string& column_name, DataTypeSerDeSPtr serde,
175
                                          bool* valid);
176
177
    Status _simdjson_write_columns_by_jsonpath(simdjson::ondemand::object* value,
178
                                               const std::vector<SlotDescriptor*>& slot_descs,
179
                                               Block& block, bool* valid);
180
    Status _append_error_msg(simdjson::ondemand::object* obj, std::string error_msg,
181
                             std::string col_name, bool* valid);
182
183
    size_t _column_index(const StringRef& name, size_t key_index);
184
185
    Status (NewJsonReader::*_vhandle_json_callback)(RuntimeState* state, Block& block,
186
                                                    const std::vector<SlotDescriptor*>& slot_descs,
187
                                                    bool* is_empty_row, bool* eof);
188
    Status _get_column_default_value(
189
            const std::vector<SlotDescriptor*>& slot_descs,
190
            const std::unordered_map<std::string, VExprContextSPtr>& col_default_value_ctx);
191
192
    Status _fill_missing_column(SlotDescriptor* slot_desc, DataTypeSerDeSPtr serde,
193
                                IColumn* column_ptr, bool* valid);
194
195
    // fe will add skip_bitmap_col to _file_slot_descs iff the target olap table has skip_bitmap_col
196
    // and the current load is a flexible partial update
197
    // flexible partial update can not be used when user specify jsonpaths, so we just fill the skip bitmap
198
    // in `_simdjson_handle_simple_json` and `_vhandle_simple_json` (which will be used when jsonpaths is not specified)
199
2.09k
    bool _should_process_skip_bitmap_col() const { return skip_bitmap_col_idx != -1; }
200
    void _append_empty_skip_bitmap_value(Block& block, size_t cur_row_count);
201
    void _set_skip_bitmap_mark(SlotDescriptor* slot_desc, IColumn* column_ptr, Block& block,
202
                               size_t cur_row_count, bool* valid);
203
    RuntimeState* _state = nullptr;
204
    RuntimeProfile* _profile = nullptr;
205
    ScannerCounter* _counter = nullptr;
206
    const TFileScanRangeParams& _params;
207
    const TFileRangeDesc& _range;
208
    io::FileSystemProperties _system_properties;
209
    io::FileDescription _file_description;
210
    const std::vector<SlotDescriptor*>& _file_slot_descs;
211
212
    io::FileReaderSPtr _file_reader;
213
    std::unique_ptr<NewPlainTextLineReader> _line_reader;
214
    bool _reader_eof;
215
    std::unique_ptr<Decompressor> _decompressor;
216
    TFileCompressType::type _file_compress_type;
217
218
    // When we fetch range doesn't start from 0 will always skip the first line
219
    bool _skip_first_line;
220
221
    std::string _line_delimiter;
222
    size_t _line_delimiter_length;
223
224
    uint32_t _next_row;
225
    size_t _total_rows;
226
227
    std::string _jsonpaths;
228
    std::string _json_root;
229
    bool _read_json_by_line;
230
    bool _strip_outer_array;
231
    bool _num_as_string;
232
    bool _fuzzy_parse;
233
234
    std::vector<std::vector<JsonPath>> _parsed_jsonpaths;
235
    std::vector<JsonPath> _parsed_json_root;
236
    bool _parsed_from_json_root = false; // to avoid parsing json root multiple times
237
238
    char _value_buffer[4 * 1024 * 1024]; // 4MB
239
    char _parse_buffer[512 * 1024];      // 512KB
240
241
    using Document = rapidjson::GenericDocument<rapidjson::UTF8<>, rapidjson::MemoryPoolAllocator<>,
242
                                                rapidjson::MemoryPoolAllocator<>>;
243
    rapidjson::MemoryPoolAllocator<> _value_allocator;
244
    rapidjson::MemoryPoolAllocator<> _parse_allocator;
245
    Document _origin_json_doc;   // origin json document object from parsed json string
246
    rapidjson::Value* _json_doc; // _json_doc equals _final_json_doc iff not set `json_root`
247
    std::unordered_map<std::string, int> _name_map;
248
249
    bool* _scanner_eof = nullptr;
250
251
    size_t _current_offset;
252
253
    io::IOContext* _io_ctx = nullptr;
254
    std::shared_ptr<io::IOContext> _io_ctx_holder;
255
256
    RuntimeProfile::Counter* _read_timer = nullptr;
257
258
    // ======SIMD JSON======
259
    // name mapping
260
    /// Hash table match `field name -> position in the block`. NOTE You can use perfect hash map.
261
    using NameMap = phmap::flat_hash_map<StringRef, size_t, StringRefHash>;
262
    NameMap _slot_desc_index;
263
    /// Cached search results for previous row (keyed as index in JSON object) - used as a hint.
264
    std::vector<NameMap::iterator> _prev_positions;
265
    /// Set of columns which already met in row. Exception is thrown if there are more than one column with the same name.
266
    std::vector<UInt8> _seen_columns;
267
    // simdjson
268
    DorisUniqueBufferPtr<uint8_t> _json_str_ptr;
269
    const uint8_t* _json_str = nullptr;
270
    static constexpr size_t _init_buffer_size = 1024 * 1024 * 8;
271
    size_t _padded_size = _init_buffer_size + simdjson::SIMDJSON_PADDING;
272
    size_t _original_doc_size = 0;
273
    std::string _simdjson_ondemand_padding_buffer;
274
    std::string _simdjson_ondemand_unscape_padding_buffer;
275
    // char _simdjson_ondemand_padding_buffer[_padded_size];
276
    simdjson::ondemand::document _original_json_doc;
277
    simdjson::ondemand::value _json_value;
278
    // for strip outer array
279
    // array_iter pointed to _array
280
    simdjson::ondemand::array_iterator _array_iter;
281
    simdjson::ondemand::array _array;
282
    std::unique_ptr<simdjson::ondemand::parser> _ondemand_json_parser;
283
    // column to default value string map
284
    std::unordered_map<std::string, std::string> _col_default_value_map;
285
286
    // From document of simdjson:
287
    // ```
288
    //   Important: a value should be consumed once. Calling get_string() twice on the same value is an error.
289
    // ```
290
    // We should cache the string_views to avoid multiple get_string() calls.
291
    struct StringViewHash {
292
0
        size_t operator()(const std::string_view& str) const {
293
0
            return std::hash<int64_t>()(reinterpret_cast<int64_t>(str.data()));
294
0
        }
295
    };
296
    struct StringViewEqual {
297
0
        bool operator()(const std::string_view& lhs, const std::string_view& rhs) const {
298
0
            return lhs.data() == rhs.data() && lhs.size() == rhs.size();
299
0
        }
300
    };
301
    std::unordered_map<std::string_view, std::string_view, StringViewHash, StringViewEqual>
302
            _cached_string_values;
303
304
    int32_t skip_bitmap_col_idx {-1};
305
306
    //Used to indicate whether it is a stream load. When loading, only data will be inserted into columnString.
307
    //If an illegal value is encountered during the load process, `_append_error_msg` should be called
308
    //instead of directly returning `Status::DataQualityError`
309
    bool _is_load = true;
310
311
    // In hive : create table xxx ROW FORMAT SERDE 'org.apache.hive.hcatalog.data.JsonSerDe';
312
    // Hive will not allow you to create columns with the same name but different case, including field names inside
313
    // structs, and will automatically convert uppercase names in create sql to lowercase.However, when Hive loads data
314
    // to table, the column names in the data may be uppercase,and there may be multiple columns with
315
    // the same name but different capitalization.We refer to the behavior of hive, convert all column names
316
    // in the data to lowercase,and use the last one as the insertion value
317
    bool _is_hive_table = false;
318
319
    // hive : org.openx.data.jsonserde.JsonSerDe, `ignore.malformed.json` prop.
320
    // If the variable is true, `null` will be inserted for llegal json format instead of return error.
321
    bool _openx_json_ignore_malformed = false;
322
323
    DataTypeSerDeSPtrs _serdes;
324
    DataTypeSerDe::FormatOptions _serde_options;
325
    // Adaptive batch size set by FileScanner.
326
    size_t _batch_size;
327
};
328
329
} // namespace doris