Coverage Report

Created: 2026-09-28 17:21

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/json/new_json_reader.cpp
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
#include "format/json/new_json_reader.h"
19
20
#include <fmt/format.h>
21
#include <gen_cpp/Metrics_types.h>
22
#include <gen_cpp/PlanNodes_types.h>
23
#include <gen_cpp/Types_types.h>
24
#include <glog/logging.h>
25
#include <rapidjson/error/en.h>
26
#include <rapidjson/reader.h>
27
#include <rapidjson/stringbuffer.h>
28
#include <rapidjson/writer.h>
29
#include <simdjson/simdjson.h> // IWYU pragma: keep
30
31
#include <algorithm>
32
#include <cinttypes>
33
#include <cstdio>
34
#include <cstring>
35
#include <map>
36
#include <memory>
37
#include <string_view>
38
#include <utility>
39
40
#include "common/compiler_util.h" // IWYU pragma: keep
41
#include "common/config.h"
42
#include "common/status.h"
43
#include "core/assert_cast.h"
44
#include "core/block/column_with_type_and_name.h"
45
#include "core/column/column.h"
46
#include "core/column/column_array.h"
47
#include "core/column/column_map.h"
48
#include "core/column/column_nullable.h"
49
#include "core/column/column_string.h"
50
#include "core/column/column_struct.h"
51
#include "core/custom_allocator.h"
52
#include "core/data_type/data_type_array.h"
53
#include "core/data_type/data_type_factory.hpp"
54
#include "core/data_type/data_type_map.h"
55
#include "core/data_type/data_type_number.h" // IWYU pragma: keep
56
#include "core/data_type/data_type_struct.h"
57
#include "core/data_type/define_primitive_type.h"
58
#include "exec/scan/scanner.h"
59
#include "exprs/json_functions.h"
60
#include "format/file_reader/new_plain_text_line_reader.h"
61
#include "io/file_factory.h"
62
#include "io/fs/buffered_reader.h"
63
#include "io/fs/file_reader.h"
64
#include "io/fs/stream_load_pipe.h"
65
#include "io/fs/tracing_file_reader.h"
66
#include "runtime/descriptors.h"
67
#include "runtime/runtime_state.h"
68
#include "util/slice.h"
69
70
namespace doris::io {
71
struct IOContext;
72
enum class FileCachePolicy : uint8_t;
73
} // namespace doris::io
74
75
namespace doris {
76
using namespace ErrorCode;
77
78
NewJsonReader::NewJsonReader(RuntimeState* state, RuntimeProfile* profile, ScannerCounter* counter,
79
                             const TFileScanRangeParams& params, const TFileRangeDesc& range,
80
                             const std::vector<SlotDescriptor*>& file_slot_descs, bool* scanner_eof,
81
                             size_t batch_size, io::IOContext* io_ctx,
82
                             std::shared_ptr<io::IOContext> io_ctx_holder)
83
1.30k
        : _vhandle_json_callback(nullptr),
84
1.30k
          _state(state),
85
1.30k
          _profile(profile),
86
1.30k
          _counter(counter),
87
1.30k
          _params(params),
88
1.30k
          _range(range),
89
1.30k
          _file_slot_descs(file_slot_descs),
90
1.30k
          _file_reader(nullptr),
91
1.30k
          _line_reader(nullptr),
92
1.30k
          _reader_eof(false),
93
1.30k
          _decompressor(nullptr),
94
1.30k
          _skip_first_line(false),
95
1.30k
          _next_row(0),
96
1.30k
          _total_rows(0),
97
1.30k
          _value_allocator(_value_buffer, sizeof(_value_buffer)),
98
1.30k
          _parse_allocator(_parse_buffer, sizeof(_parse_buffer)),
99
1.30k
          _origin_json_doc(&_value_allocator, sizeof(_parse_buffer), &_parse_allocator),
100
1.30k
          _scanner_eof(scanner_eof),
101
1.30k
          _current_offset(0),
102
1.30k
          _io_ctx(io_ctx),
103
1.30k
          _io_ctx_holder(std::move(io_ctx_holder)),
104
1.30k
          _batch_size(std::max(batch_size, 1UL)) {
105
1.30k
    if (_io_ctx == nullptr && _io_ctx_holder) {
106
0
        _io_ctx = _io_ctx_holder.get();
107
0
    }
108
1.30k
    _read_timer = ADD_TIMER(_profile, "ReadTime");
109
1.30k
    if (_range.__isset.compress_type) {
110
        // for compatibility
111
0
        _file_compress_type = _range.compress_type;
112
1.30k
    } else {
113
1.30k
        _file_compress_type = _params.compress_type;
114
1.30k
    }
115
1.30k
    _init_system_properties();
116
1.30k
    _init_file_description();
117
1.30k
}
118
119
NewJsonReader::NewJsonReader(RuntimeProfile* profile, const TFileScanRangeParams& params,
120
                             const TFileRangeDesc& range,
121
                             const std::vector<SlotDescriptor*>& file_slot_descs, size_t batch_size,
122
                             io::IOContext* io_ctx, std::shared_ptr<io::IOContext> io_ctx_holder)
123
2
        : _vhandle_json_callback(nullptr),
124
2
          _state(nullptr),
125
2
          _profile(profile),
126
2
          _params(params),
127
2
          _range(range),
128
2
          _file_slot_descs(file_slot_descs),
129
2
          _line_reader(nullptr),
130
2
          _reader_eof(false),
131
2
          _decompressor(nullptr),
132
2
          _skip_first_line(false),
133
2
          _next_row(0),
134
2
          _total_rows(0),
135
2
          _value_allocator(_value_buffer, sizeof(_value_buffer)),
136
2
          _parse_allocator(_parse_buffer, sizeof(_parse_buffer)),
137
2
          _origin_json_doc(&_value_allocator, sizeof(_parse_buffer), &_parse_allocator),
138
2
          _io_ctx(io_ctx),
139
2
          _io_ctx_holder(std::move(io_ctx_holder)),
140
2
          _batch_size(std::max(batch_size, 1UL)) {
141
2
    if (_io_ctx == nullptr && _io_ctx_holder) {
142
0
        _io_ctx = _io_ctx_holder.get();
143
0
    }
144
2
    if (_range.__isset.compress_type) {
145
        // for compatibility
146
0
        _file_compress_type = _range.compress_type;
147
2
    } else {
148
2
        _file_compress_type = _params.compress_type;
149
2
    }
150
2
    _init_system_properties();
151
2
    _init_file_description();
152
2
}
153
154
1.30k
NewJsonReader::~NewJsonReader() = default;
155
156
1.30k
void NewJsonReader::_init_system_properties() {
157
1.30k
    if (_range.__isset.file_type) {
158
        // for compatibility
159
0
        _system_properties.system_type = _range.file_type;
160
1.30k
    } else {
161
1.30k
        _system_properties.system_type = _params.file_type;
162
1.30k
    }
163
1.30k
    _system_properties.properties = _params.properties;
164
1.30k
    _system_properties.hdfs_params = _params.hdfs_params;
165
1.30k
    if (_params.__isset.broker_addresses) {
166
0
        _system_properties.broker_addresses.assign(_params.broker_addresses.begin(),
167
0
                                                   _params.broker_addresses.end());
168
0
    }
169
1.30k
}
170
171
1.30k
void NewJsonReader::_init_file_description() {
172
1.30k
    _file_description.path = _range.path;
173
1.30k
    _file_description.file_size = _range.__isset.file_size ? _range.file_size : -1;
174
175
1.30k
    if (_range.__isset.fs_name) {
176
0
        _file_description.fs_name = _range.fs_name;
177
0
    }
178
1.30k
    if (_range.__isset.file_cache_admission) {
179
0
        _file_description.file_cache_admission = _range.file_cache_admission;
180
0
    }
181
1.30k
}
182
183
Status NewJsonReader::init_reader(
184
        const std::unordered_map<std::string, VExprContextSPtr>& col_default_value_ctx,
185
1.30k
        bool is_load) {
186
1.30k
    _is_load = is_load;
187
188
    // generate _col_default_value_map
189
1.30k
    RETURN_IF_ERROR(_get_column_default_value(_file_slot_descs, col_default_value_ctx));
190
191
    //use serde insert data to column.
192
1.30k
    for (auto* slot_desc : _file_slot_descs) {
193
1.30k
        _serdes.emplace_back(slot_desc->get_data_type_ptr()->get_serde());
194
1.30k
    }
195
196
    // create decompressor.
197
    // _decompressor may be nullptr if this is not a compressed file
198
1.30k
    RETURN_IF_ERROR(Decompressor::create_decompressor(_file_compress_type, &_decompressor));
199
200
1.30k
    RETURN_IF_ERROR(_simdjson_init_reader());
201
1.30k
    return Status::OK();
202
1.30k
}
203
204
// ---- Unified init_reader(ReaderInitContext*) overrides ----
205
206
0
Status NewJsonReader::_open_file_reader(ReaderInitContext* /*ctx*/) {
207
0
    RETURN_IF_ERROR(_get_range_params());
208
0
    RETURN_IF_ERROR(_open_file_reader(false));
209
0
    return Status::OK();
210
0
}
211
212
0
Status NewJsonReader::_do_init_reader(ReaderInitContext* base_ctx) {
213
0
    auto* ctx = checked_context_cast<JsonInitContext>(base_ctx);
214
0
    _is_load = ctx->is_load;
215
216
0
    RETURN_IF_ERROR(_get_column_default_value(_file_slot_descs, *ctx->col_default_value_ctx));
217
0
    for (auto* slot_desc : _file_slot_descs) {
218
0
        _serdes.emplace_back(slot_desc->get_data_type_ptr()->get_serde());
219
0
    }
220
221
    // Create decompressor (needed by line reader below)
222
0
    RETURN_IF_ERROR(Decompressor::create_decompressor(_file_compress_type, &_decompressor));
223
224
0
    if (LIKELY(_read_json_by_line)) {
225
0
        RETURN_IF_ERROR(_open_line_reader());
226
0
    }
227
0
    RETURN_IF_ERROR(_parse_jsonpath_and_json_root());
228
229
0
    if (_parsed_jsonpaths.empty()) {
230
0
        _vhandle_json_callback = &NewJsonReader::_simdjson_handle_simple_json;
231
0
    } else {
232
0
        if (_strip_outer_array) {
233
0
            _vhandle_json_callback = &NewJsonReader::_simdjson_handle_flat_array_complex_json;
234
0
        } else {
235
0
            _vhandle_json_callback = &NewJsonReader::_simdjson_handle_nested_complex_json;
236
0
        }
237
0
    }
238
0
    _ondemand_json_parser = std::make_unique<simdjson::ondemand::parser>();
239
0
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
240
0
        _slot_desc_index[StringRef {_file_slot_descs[i]->col_name()}] = i;
241
0
        if (_file_slot_descs[i]->is_skip_bitmap_col()) {
242
0
            skip_bitmap_col_idx = i;
243
0
        }
244
0
    }
245
0
    _simdjson_ondemand_padding_buffer.resize(_padded_size);
246
0
    _simdjson_ondemand_unscape_padding_buffer.resize(_padded_size);
247
0
    return Status::OK();
248
0
}
249
250
5
void NewJsonReader::set_batch_size(size_t batch_size) {
251
    // 0 means "not set" / "use default" for the row-based readers; we must
252
    // never let _batch_size be 0 because _do_get_next_block uses it as the
253
    // upper bound of a `while (block->rows() < batch_size)` loop and a 0
254
    // would make the reader return without setting eof, causing the scanner
255
    // to spin on empty blocks.
256
5
    _batch_size = std::max(batch_size, 1UL);
257
5
}
258
259
3.36k
Status NewJsonReader::_do_get_next_block(Block* block, size_t* read_rows, bool* eof) {
260
3.36k
    if (_reader_eof) {
261
1.30k
        *eof = true;
262
1.30k
        return Status::OK();
263
1.30k
    }
264
265
2.05k
    const auto batch_size = _batch_size;
266
2.05k
    const auto max_block_bytes = _state->preferred_block_size_bytes();
267
268
38.9k
    while (block->rows() < batch_size && !_reader_eof && (block->bytes() < max_block_bytes)) {
269
36.9k
        if (UNLIKELY(_read_json_by_line && _skip_first_line)) {
270
736
            RETURN_IF_ERROR(_line_reader->skip_split_prefix(_range.start_offset, _line_delimiter,
271
736
                                                            &_reader_eof, _io_ctx));
272
736
            _skip_first_line = false;
273
736
            continue;
274
736
        }
275
276
36.1k
        bool is_empty_row = false;
277
278
36.1k
        RETURN_IF_ERROR(
279
36.1k
                _read_json_column(_state, *block, _file_slot_descs, &is_empty_row, &_reader_eof));
280
36.1k
        if (is_empty_row) {
281
            // Read empty row, just continue
282
34.0k
            continue;
283
34.0k
        }
284
2.09k
        ++(*read_rows);
285
2.09k
    }
286
287
2.05k
    return Status::OK();
288
2.05k
}
289
290
Status NewJsonReader::_get_columns_impl(
291
0
        std::unordered_map<std::string, DataTypePtr>* name_to_type) {
292
0
    for (const auto& slot : _file_slot_descs) {
293
0
        name_to_type->emplace(slot->col_name(), slot->type());
294
0
    }
295
0
    return Status::OK();
296
0
}
297
298
// init decompressor, file reader and line reader for parsing schema
299
0
Status NewJsonReader::init_schema_reader() {
300
0
    RETURN_IF_ERROR(_get_range_params());
301
    // create decompressor.
302
    // _decompressor may be nullptr if this is not a compressed file
303
0
    RETURN_IF_ERROR(Decompressor::create_decompressor(_file_compress_type, &_decompressor));
304
0
    RETURN_IF_ERROR(_open_file_reader(true));
305
0
    if (_read_json_by_line) {
306
0
        RETURN_IF_ERROR(_open_line_reader());
307
0
    }
308
    // generate _parsed_jsonpaths and _parsed_json_root
309
0
    RETURN_IF_ERROR(_parse_jsonpath_and_json_root());
310
0
    return Status::OK();
311
0
}
312
313
Status NewJsonReader::get_parsed_schema(std::vector<std::string>* col_names,
314
0
                                        std::vector<DataTypePtr>* col_types) {
315
0
    bool eof = false;
316
0
    const uint8_t* json_str = nullptr;
317
0
    DorisUniqueBufferPtr<uint8_t> json_str_ptr;
318
0
    size_t size = 0;
319
0
    if (_line_reader != nullptr) {
320
0
        RETURN_IF_ERROR(_line_reader->read_line(&json_str, &size, &eof, _io_ctx));
321
0
    } else {
322
0
        size_t read_size = 0;
323
0
        RETURN_IF_ERROR(_read_one_message(&json_str_ptr, &read_size));
324
0
        json_str = json_str_ptr.get();
325
0
        size = read_size;
326
0
        if (read_size == 0) {
327
0
            eof = true;
328
0
        }
329
0
    }
330
331
0
    if (size == 0 || eof) {
332
0
        return Status::EndOfFile("Empty file.");
333
0
    }
334
335
    // clear memory here.
336
0
    _value_allocator.Clear();
337
0
    _parse_allocator.Clear();
338
0
    bool has_parse_error = false;
339
340
    // parse jsondata to JsonDoc
341
    // As the issue: https://github.com/Tencent/rapidjson/issues/1458
342
    // Now, rapidjson only support uint64_t, So lagreint load cause bug. We use kParseNumbersAsStringsFlag.
343
0
    if (_num_as_string) {
344
0
        has_parse_error =
345
0
                _origin_json_doc.Parse<rapidjson::kParseNumbersAsStringsFlag>((char*)json_str, size)
346
0
                        .HasParseError();
347
0
    } else {
348
0
        has_parse_error = _origin_json_doc.Parse((char*)json_str, size).HasParseError();
349
0
    }
350
351
0
    if (has_parse_error) {
352
0
        return Status::DataQualityError(
353
0
                "Parse json data for JsonDoc failed. code: {}, error info: {}",
354
0
                _origin_json_doc.GetParseError(),
355
0
                rapidjson::GetParseError_En(_origin_json_doc.GetParseError()));
356
0
    }
357
358
    // set json root
359
0
    if (!_parsed_json_root.empty()) {
360
0
        _json_doc = JsonFunctions::get_json_object_from_parsed_json(
361
0
                _parsed_json_root, &_origin_json_doc, _origin_json_doc.GetAllocator());
362
0
        if (_json_doc == nullptr) {
363
0
            return Status::DataQualityError("JSON Root not found.");
364
0
        }
365
0
    } else {
366
0
        _json_doc = &_origin_json_doc;
367
0
    }
368
369
0
    if (_json_doc->IsArray() && !_strip_outer_array) {
370
0
        return Status::DataQualityError(
371
0
                "JSON data is array-object, `strip_outer_array` must be TRUE.");
372
0
    }
373
0
    if (!_json_doc->IsArray() && _strip_outer_array) {
374
0
        return Status::DataQualityError(
375
0
                "JSON data is not an array-object, `strip_outer_array` must be FALSE.");
376
0
    }
377
378
0
    rapidjson::Value* objectValue = nullptr;
379
0
    if (_json_doc->IsArray()) {
380
0
        if (_json_doc->Size() == 0) {
381
            // may be passing an empty json, such as "[]"
382
0
            return Status::InternalError<false>("Empty first json line");
383
0
        }
384
0
        objectValue = &(*_json_doc)[0];
385
0
    } else {
386
0
        objectValue = _json_doc;
387
0
    }
388
389
0
    if (!objectValue->IsObject()) {
390
0
        return Status::DataQualityError("JSON data is not an object. but: {}",
391
0
                                        objectValue->GetType());
392
0
    }
393
394
    // use jsonpaths to col_names
395
0
    if (!_parsed_jsonpaths.empty()) {
396
0
        for (auto& _parsed_jsonpath : _parsed_jsonpaths) {
397
0
            size_t len = _parsed_jsonpath.size();
398
0
            if (len == 0) {
399
0
                return Status::InvalidArgument("It's invalid jsonpaths.");
400
0
            }
401
0
            std::string key = _parsed_jsonpath[len - 1].key;
402
0
            col_names->emplace_back(key);
403
0
            col_types->emplace_back(
404
0
                    DataTypeFactory::instance().create_data_type(PrimitiveType::TYPE_STRING, true));
405
0
        }
406
0
        return Status::OK();
407
0
    }
408
409
0
    for (int i = 0; i < objectValue->MemberCount(); ++i) {
410
0
        auto it = objectValue->MemberBegin() + i;
411
0
        col_names->emplace_back(it->name.GetString());
412
0
        col_types->emplace_back(make_nullable(std::make_shared<DataTypeString>()));
413
0
    }
414
0
    return Status::OK();
415
0
}
416
417
1.30k
Status NewJsonReader::_get_range_params() {
418
1.30k
    if (!_params.__isset.file_attributes) {
419
0
        return Status::InternalError<false>("BE cat get file_attributes");
420
0
    }
421
422
    // get line_delimiter
423
1.30k
    if (_params.file_attributes.__isset.text_params &&
424
1.30k
        _params.file_attributes.text_params.__isset.line_delimiter) {
425
1.30k
        _line_delimiter = _params.file_attributes.text_params.line_delimiter;
426
1.30k
        _line_delimiter_length = _line_delimiter.size();
427
1.30k
    }
428
429
1.30k
    if (_params.file_attributes.__isset.jsonpaths) {
430
0
        _jsonpaths = _params.file_attributes.jsonpaths;
431
0
    }
432
1.30k
    if (_params.file_attributes.__isset.json_root) {
433
0
        _json_root = _params.file_attributes.json_root;
434
0
    }
435
1.30k
    if (_params.file_attributes.__isset.read_json_by_line) {
436
1.30k
        _read_json_by_line = _params.file_attributes.read_json_by_line;
437
1.30k
    }
438
1.30k
    if (_params.file_attributes.__isset.strip_outer_array) {
439
1.30k
        _strip_outer_array = _params.file_attributes.strip_outer_array;
440
1.30k
    }
441
1.30k
    if (_params.file_attributes.__isset.num_as_string) {
442
1.30k
        _num_as_string = _params.file_attributes.num_as_string;
443
1.30k
    }
444
1.30k
    if (_params.file_attributes.__isset.fuzzy_parse) {
445
1.30k
        _fuzzy_parse = _params.file_attributes.fuzzy_parse;
446
1.30k
    }
447
1.30k
    if (_range.table_format_params.table_format_type == "hive") {
448
0
        _is_hive_table = true;
449
0
    }
450
1.30k
    if (_params.file_attributes.__isset.openx_json_ignore_malformed) {
451
1.30k
        _openx_json_ignore_malformed = _params.file_attributes.openx_json_ignore_malformed;
452
1.30k
    }
453
1.30k
    return Status::OK();
454
1.30k
}
455
456
1
Status json_reader_detail::append_null_for_malformed_json(Block& block) {
457
2
    for (int i = 0; i < block.columns(); ++i) {
458
1
        auto& column_with_type = block.get_by_position(i);
459
1
        if (!is_column_nullable(*column_with_type.column)) [[unlikely]] {
460
0
            return Status::DataQualityError("malformed json, but the column `{}` is not nullable.",
461
0
                                            column_with_type.column->get_name());
462
0
        }
463
1
        auto column = IColumn::mutate(std::move(column_with_type.column));
464
1
        assert_cast<ColumnNullable*>(column.get())->insert_default();
465
1
        column_with_type.column = std::move(column);
466
1
    }
467
1
    return Status::OK();
468
1
}
469
470
1
void json_reader_detail::truncate_block_to_rows(Block& block, size_t num_rows) {
471
2
    for (int i = 0; i < block.columns(); ++i) {
472
1
        auto& column_with_type = block.get_by_position(i);
473
1
        auto column = IColumn::mutate(std::move(column_with_type.column));
474
1
        if (column->size() > num_rows) {
475
1
            column->pop_back(column->size() - num_rows);
476
1
        }
477
1
        column_with_type.column = std::move(column);
478
1
    }
479
1
}
480
481
1
void json_reader_detail::pop_back_last_inserted_value(Block& block, size_t column_index) {
482
1
    auto& column = block.get_by_position(column_index).column;
483
1
    auto mutable_column = IColumn::mutate(std::move(column));
484
1
    mutable_column->pop_back(1);
485
1
    column = std::move(mutable_column);
486
1
}
487
488
1.30k
Status NewJsonReader::_open_file_reader(bool need_schema) {
489
1.30k
    int64_t start_offset = _range.start_offset;
490
1.30k
    if (start_offset != 0) {
491
        // Include the whole delimiter when the split starts inside it, so skipping the first
492
        // partial line cannot discard the next complete JSON record.
493
736
        start_offset -= std::min<int64_t>(start_offset, _line_delimiter_length);
494
736
    }
495
496
1.30k
    _current_offset = start_offset;
497
498
1.30k
    if (_params.file_type == TFileType::FILE_STREAM) {
499
        // Due to http_stream needs to pre read a portion of the data to parse column information, so it is set to true here
500
0
        RETURN_IF_ERROR(FileFactory::create_pipe_reader(_range.load_id, &_file_reader, _state,
501
0
                                                        need_schema));
502
1.30k
    } else {
503
1.30k
        _file_description.mtime = _range.__isset.modification_time ? _range.modification_time : 0;
504
1.30k
        io::FileReaderOptions reader_options = FileFactory::get_reader_options(
505
1.30k
                _state ? _state->query_options() : _default_query_options, _file_description);
506
1.30k
        io::FileReaderSPtr file_reader;
507
1.30k
        if (_io_ctx_holder) {
508
0
            file_reader = DORIS_TRY(io::DelegateReader::create_file_reader(
509
0
                    _profile, _system_properties, _file_description, reader_options,
510
0
                    io::DelegateReader::AccessMode::SEQUENTIAL,
511
0
                    std::static_pointer_cast<const io::IOContext>(_io_ctx_holder),
512
0
                    io::PrefetchRange(_range.start_offset, _range.size)));
513
1.30k
        } else {
514
1.30k
            file_reader = DORIS_TRY(io::DelegateReader::create_file_reader(
515
1.30k
                    _profile, _system_properties, _file_description, reader_options,
516
1.30k
                    io::DelegateReader::AccessMode::SEQUENTIAL, _io_ctx,
517
1.30k
                    io::PrefetchRange(_range.start_offset, _range.size)));
518
1.30k
        }
519
1.30k
        _file_reader = _io_ctx && _io_ctx->file_reader_stats
520
1.30k
                               ? std::make_shared<io::TracingFileReader>(std::move(file_reader),
521
0
                                                                         _io_ctx->file_reader_stats)
522
1.30k
                               : file_reader;
523
1.30k
    }
524
1.30k
    return Status::OK();
525
1.30k
}
526
527
1.30k
Status NewJsonReader::_open_line_reader() {
528
1.30k
    int64_t size = _range.size;
529
1.30k
    if (_range.start_offset != 0) {
530
        // Preserve the original range end after moving the start backwards.
531
736
        size += _range.start_offset - _current_offset;
532
736
        _skip_first_line = true;
533
736
    } else {
534
565
        _skip_first_line = false;
535
565
    }
536
1.30k
    _line_reader = NewPlainTextLineReader::create_unique(
537
1.30k
            _profile, _file_reader, _decompressor.get(),
538
1.30k
            std::make_shared<PlainTextLineReaderCtx>(_line_delimiter, _line_delimiter_length,
539
1.30k
                                                     false),
540
1.30k
            size, _current_offset);
541
1.30k
    return Status::OK();
542
1.30k
}
543
544
1.30k
Status NewJsonReader::_parse_jsonpath_and_json_root() {
545
    // parse jsonpaths
546
1.30k
    if (!_jsonpaths.empty()) {
547
0
        rapidjson::Document jsonpaths_doc;
548
0
        if (!jsonpaths_doc.Parse(_jsonpaths.c_str(), _jsonpaths.length()).HasParseError()) {
549
0
            if (!jsonpaths_doc.IsArray()) {
550
0
                return Status::InvalidJsonPath("Invalid json path: {}", _jsonpaths);
551
0
            }
552
0
            for (int i = 0; i < jsonpaths_doc.Size(); i++) {
553
0
                const rapidjson::Value& path = jsonpaths_doc[i];
554
0
                if (!path.IsString()) {
555
0
                    return Status::InvalidJsonPath("Invalid json path: {}", _jsonpaths);
556
0
                }
557
0
                std::string json_path = path.GetString();
558
                // $ -> $. in json_path
559
0
                if (UNLIKELY(json_path.size() == 1 && json_path[0] == '$')) {
560
0
                    json_path.insert(1, ".");
561
0
                }
562
0
                std::vector<JsonPath> parsed_paths;
563
0
                JsonFunctions::parse_json_paths(json_path, &parsed_paths);
564
0
                _parsed_jsonpaths.push_back(std::move(parsed_paths));
565
0
            }
566
567
0
        } else {
568
0
            return Status::InvalidJsonPath("Invalid json path: {}", _jsonpaths);
569
0
        }
570
0
    }
571
572
    // parse jsonroot
573
1.30k
    if (!_json_root.empty()) {
574
0
        std::string json_root = _json_root;
575
        //  $ -> $. in json_root
576
0
        if (json_root.size() == 1 && json_root[0] == '$') {
577
0
            json_root.insert(1, ".");
578
0
        }
579
0
        JsonFunctions::parse_json_paths(json_root, &_parsed_json_root);
580
0
    }
581
1.30k
    return Status::OK();
582
1.30k
}
583
584
Status NewJsonReader::_read_json_column(RuntimeState* state, Block& block,
585
                                        const std::vector<SlotDescriptor*>& slot_descs,
586
36.1k
                                        bool* is_empty_row, bool* eof) {
587
36.1k
    return (this->*_vhandle_json_callback)(state, block, slot_descs, is_empty_row, eof);
588
36.1k
}
589
590
Status NewJsonReader::_read_one_message(DorisUniqueBufferPtr<uint8_t>* file_buf,
591
0
                                        size_t* read_size) {
592
0
    switch (_params.file_type) {
593
0
    case TFileType::FILE_LOCAL:
594
0
        [[fallthrough]];
595
0
    case TFileType::FILE_HDFS:
596
0
    case TFileType::FILE_HTTP:
597
0
        [[fallthrough]];
598
0
    case TFileType::FILE_S3: {
599
0
        size_t file_size = _file_reader->size();
600
0
        *file_buf = make_unique_buffer<uint8_t>(file_size);
601
0
        Slice result(file_buf->get(), file_size);
602
0
        RETURN_IF_ERROR(_file_reader->read_at(_current_offset, result, read_size, _io_ctx));
603
0
        _current_offset += *read_size;
604
0
        break;
605
0
    }
606
0
    case TFileType::FILE_STREAM: {
607
0
        RETURN_IF_ERROR(_read_one_message_from_pipe(file_buf, read_size));
608
0
        break;
609
0
    }
610
0
    default: {
611
0
        return Status::NotSupported<false>("no supported file reader type: {}", _params.file_type);
612
0
    }
613
0
    }
614
0
    return Status::OK();
615
0
}
616
617
Status NewJsonReader::_read_one_message_from_pipe(DorisUniqueBufferPtr<uint8_t>* file_buf,
618
0
                                                  size_t* read_size) {
619
0
    auto* stream_load_pipe = dynamic_cast<io::StreamLoadPipe*>(_file_reader.get());
620
621
    // first read: read from the pipe once.
622
0
    RETURN_IF_ERROR(stream_load_pipe->read_one_message(file_buf, read_size));
623
624
    // When the file is not chunked, the entire file has already been read.
625
0
    if (!stream_load_pipe->is_chunked_transfer()) {
626
0
        return Status::OK();
627
0
    }
628
629
0
    std::vector<uint8_t> buf;
630
0
    uint64_t cur_size = 0;
631
632
    // second read: continuously read data from the pipe until all data is read.
633
0
    DorisUniqueBufferPtr<uint8_t> read_buf;
634
0
    size_t read_buf_size = 0;
635
0
    while (true) {
636
0
        RETURN_IF_ERROR(stream_load_pipe->read_one_message(&read_buf, &read_buf_size));
637
0
        if (read_buf_size == 0) {
638
0
            break;
639
0
        } else {
640
0
            buf.insert(buf.end(), read_buf.get(), read_buf.get() + read_buf_size);
641
0
            cur_size += read_buf_size;
642
0
            read_buf_size = 0;
643
0
            read_buf.reset();
644
0
        }
645
0
    }
646
647
    // No data is available during the second read.
648
0
    if (cur_size == 0) {
649
0
        return Status::OK();
650
0
    }
651
652
0
    DorisUniqueBufferPtr<uint8_t> total_buf = make_unique_buffer<uint8_t>(cur_size + *read_size);
653
654
    // copy the data during the first read
655
0
    memcpy(total_buf.get(), file_buf->get(), *read_size);
656
657
    // copy the data during the second read
658
0
    memcpy(total_buf.get() + *read_size, buf.data(), cur_size);
659
0
    *file_buf = std::move(total_buf);
660
0
    *read_size += cur_size;
661
0
    return Status::OK();
662
0
}
663
664
// ---------SIMDJSON----------
665
// simdjson, replace none simdjson function if it is ready
666
1.30k
Status NewJsonReader::_simdjson_init_reader() {
667
1.30k
    RETURN_IF_ERROR(_get_range_params());
668
669
1.30k
    RETURN_IF_ERROR(_open_file_reader(false));
670
1.30k
    if (LIKELY(_read_json_by_line)) {
671
1.30k
        RETURN_IF_ERROR(_open_line_reader());
672
1.30k
    }
673
674
    // generate _parsed_jsonpaths and _parsed_json_root
675
1.30k
    RETURN_IF_ERROR(_parse_jsonpath_and_json_root());
676
677
    //improve performance
678
1.30k
    if (_parsed_jsonpaths.empty()) { // input is a simple json-string
679
1.30k
        _vhandle_json_callback = &NewJsonReader::_simdjson_handle_simple_json;
680
1.30k
    } else { // input is a complex json-string and a json-path
681
0
        if (_strip_outer_array) {
682
0
            _vhandle_json_callback = &NewJsonReader::_simdjson_handle_flat_array_complex_json;
683
0
        } else {
684
0
            _vhandle_json_callback = &NewJsonReader::_simdjson_handle_nested_complex_json;
685
0
        }
686
0
    }
687
1.30k
    _ondemand_json_parser = std::make_unique<simdjson::ondemand::parser>();
688
2.60k
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
689
1.30k
        _slot_desc_index[StringRef {_file_slot_descs[i]->col_name()}] = i;
690
1.30k
        if (_file_slot_descs[i]->is_skip_bitmap_col()) {
691
0
            skip_bitmap_col_idx = i;
692
0
        }
693
1.30k
    }
694
1.30k
    _simdjson_ondemand_padding_buffer.resize(_padded_size);
695
1.30k
    _simdjson_ondemand_unscape_padding_buffer.resize(_padded_size);
696
1.30k
    return Status::OK();
697
1.30k
}
698
699
Status NewJsonReader::_handle_simdjson_error(simdjson::simdjson_error& error, Block& block,
700
0
                                             size_t num_rows, bool* eof) {
701
0
    fmt::memory_buffer error_msg;
702
0
    fmt::format_to(error_msg, "Parse json data failed. code: {}, error info: {}", error.error(),
703
0
                   error.what());
704
0
    _counter->num_rows_filtered++;
705
    // Before continuing to process other rows, we need to first clean the fail parsed row.
706
0
    json_reader_detail::truncate_block_to_rows(block, num_rows);
707
708
0
    RETURN_IF_ERROR(_state->append_error_msg_to_file(
709
0
            [&]() -> std::string {
710
0
                return std::string(_simdjson_ondemand_padding_buffer.data(), _original_doc_size);
711
0
            },
712
0
            [&]() -> std::string { return fmt::to_string(error_msg); }));
713
0
    return Status::OK();
714
0
}
715
716
Status NewJsonReader::_simdjson_handle_simple_json(RuntimeState* /*state*/, Block& block,
717
                                                   const std::vector<SlotDescriptor*>& slot_descs,
718
36.1k
                                                   bool* is_empty_row, bool* eof) {
719
    // simple json
720
36.1k
    size_t size = 0;
721
36.1k
    simdjson::error_code error;
722
36.1k
    size_t num_rows = block.rows();
723
36.1k
    try {
724
        // step1: get and parse buf to get json doc
725
36.1k
        RETURN_IF_ERROR(_simdjson_parse_json(&size, is_empty_row, eof, &error));
726
36.1k
        if (size == 0 || *eof) {
727
34.0k
            *is_empty_row = true;
728
34.0k
            return Status::OK();
729
34.0k
        }
730
731
        // step2: get json value by json doc
732
2.09k
        Status st = _get_json_value(&size, eof, &error, is_empty_row);
733
2.09k
        if (st.is<DATA_QUALITY_ERROR>()) {
734
0
            if (_is_load) {
735
0
                return Status::OK();
736
0
            } else if (_openx_json_ignore_malformed) {
737
0
                RETURN_IF_ERROR(json_reader_detail::append_null_for_malformed_json(block));
738
0
                return Status::OK();
739
0
            }
740
0
        }
741
742
2.09k
        RETURN_IF_ERROR(st);
743
2.09k
        if (*is_empty_row || *eof) {
744
0
            return Status::OK();
745
0
        }
746
747
        // step 3: write columns by json value
748
2.09k
        RETURN_IF_ERROR(
749
2.09k
                _simdjson_handle_simple_json_write_columns(block, slot_descs, is_empty_row, eof));
750
2.09k
    } catch (simdjson::simdjson_error& e) {
751
0
        RETURN_IF_ERROR(_handle_simdjson_error(e, block, num_rows, eof));
752
0
        if (*_scanner_eof) {
753
            // When _scanner_eof is true and valid is false, it means that we have encountered
754
            // unqualified data and decided to stop the scan.
755
0
            *is_empty_row = true;
756
0
            return Status::OK();
757
0
        }
758
0
    }
759
760
2.09k
    return Status::OK();
761
36.1k
}
762
763
Status NewJsonReader::_simdjson_handle_simple_json_write_columns(
764
        Block& block, const std::vector<SlotDescriptor*>& slot_descs, bool* is_empty_row,
765
2.09k
        bool* eof) {
766
2.09k
    simdjson::ondemand::object objectValue;
767
2.09k
    size_t num_rows = block.rows();
768
2.09k
    bool valid = false;
769
2.09k
    try {
770
2.09k
        if (_json_value.type() == simdjson::ondemand::json_type::array) {
771
0
            _array = _json_value.get_array();
772
0
            if (_array.count_elements() == 0) {
773
                // may be passing an empty json, such as "[]"
774
0
                RETURN_IF_ERROR(_append_error_msg(nullptr, "Empty json line", "", nullptr));
775
0
                if (*_scanner_eof) {
776
0
                    *is_empty_row = true;
777
0
                    return Status::OK();
778
0
                }
779
0
                return Status::OK();
780
0
            }
781
782
0
            _array_iter = _array.begin();
783
0
            while (true) {
784
0
                objectValue = *_array_iter;
785
0
                RETURN_IF_ERROR(
786
0
                        _simdjson_set_column_value(&objectValue, block, slot_descs, &valid));
787
0
                if (!valid) {
788
0
                    if (*_scanner_eof) {
789
                        // When _scanner_eof is true and valid is false, it means that we have encountered
790
                        // unqualified data and decided to stop the scan.
791
0
                        *is_empty_row = true;
792
0
                        return Status::OK();
793
0
                    }
794
0
                }
795
0
                ++_array_iter;
796
0
                if (_array_iter == _array.end()) {
797
                    // Hint to read next json doc
798
0
                    break;
799
0
                }
800
0
            }
801
2.09k
        } else {
802
2.09k
            objectValue = _json_value;
803
2.09k
            RETURN_IF_ERROR(_simdjson_set_column_value(&objectValue, block, slot_descs, &valid));
804
2.09k
            if (!valid) {
805
0
                if (*_scanner_eof) {
806
0
                    *is_empty_row = true;
807
0
                    return Status::OK();
808
0
                }
809
0
            }
810
2.09k
            *is_empty_row = false;
811
2.09k
        }
812
2.09k
    } catch (simdjson::simdjson_error& e) {
813
0
        RETURN_IF_ERROR(_handle_simdjson_error(e, block, num_rows, eof));
814
0
        if (!valid) {
815
0
            if (*_scanner_eof) {
816
0
                *is_empty_row = true;
817
0
                return Status::OK();
818
0
            }
819
0
        }
820
0
    }
821
2.09k
    return Status::OK();
822
2.09k
}
823
824
Status NewJsonReader::_simdjson_handle_flat_array_complex_json(
825
        RuntimeState* /*state*/, Block& block, const std::vector<SlotDescriptor*>& slot_descs,
826
0
        bool* is_empty_row, bool* eof) {
827
    // array complex json
828
0
    size_t size = 0;
829
0
    simdjson::error_code error;
830
0
    size_t num_rows = block.rows();
831
0
    try {
832
        // step1: get and parse buf to get json doc
833
0
        RETURN_IF_ERROR(_simdjson_parse_json(&size, is_empty_row, eof, &error));
834
0
        if (size == 0 || *eof) {
835
0
            *is_empty_row = true;
836
0
            return Status::OK();
837
0
        }
838
839
        // step2: get json value by json doc
840
0
        Status st = _get_json_value(&size, eof, &error, is_empty_row);
841
0
        if (st.is<DATA_QUALITY_ERROR>()) {
842
0
            return Status::OK();
843
0
        }
844
0
        RETURN_IF_ERROR(st);
845
0
        if (*is_empty_row) {
846
0
            return Status::OK();
847
0
        }
848
849
        // step 3: write columns by json value
850
0
        RETURN_IF_ERROR(_simdjson_handle_flat_array_complex_json_write_columns(block, slot_descs,
851
0
                                                                               is_empty_row, eof));
852
0
    } catch (simdjson::simdjson_error& e) {
853
0
        RETURN_IF_ERROR(_handle_simdjson_error(e, block, num_rows, eof));
854
0
        if (*_scanner_eof) {
855
            // When _scanner_eof is true and valid is false, it means that we have encountered
856
            // unqualified data and decided to stop the scan.
857
0
            *is_empty_row = true;
858
0
            return Status::OK();
859
0
        }
860
0
    }
861
862
0
    return Status::OK();
863
0
}
864
865
Status NewJsonReader::_simdjson_handle_flat_array_complex_json_write_columns(
866
        Block& block, const std::vector<SlotDescriptor*>& slot_descs, bool* is_empty_row,
867
0
        bool* eof) {
868
// Advance one row in array list, if it is the endpoint, stop advance and break the loop
869
0
#define ADVANCE_ROW()                  \
870
0
    ++_array_iter;                     \
871
0
    if (_array_iter == _array.end()) { \
872
0
        break;                         \
873
0
    }
874
875
0
    simdjson::ondemand::object cur;
876
0
    size_t num_rows = block.rows();
877
0
    try {
878
0
        bool valid = true;
879
0
        _array = _json_value.get_array();
880
0
        _array_iter = _array.begin();
881
882
0
        while (true) {
883
0
            cur = (*_array_iter).get_object();
884
            // extract root
885
0
            if (!_parsed_from_json_root && !_parsed_json_root.empty()) {
886
0
                simdjson::ondemand::value val;
887
0
                Status st = JsonFunctions::extract_from_object(cur, _parsed_json_root, &val);
888
0
                if (UNLIKELY(!st.ok())) {
889
0
                    if (st.is<NOT_FOUND>()) {
890
0
                        RETURN_IF_ERROR(_append_error_msg(nullptr, st.to_string(), "", nullptr));
891
0
                        ADVANCE_ROW();
892
0
                        continue;
893
0
                    }
894
0
                    return st;
895
0
                }
896
0
                if (val.type() != simdjson::ondemand::json_type::object) {
897
0
                    RETURN_IF_ERROR(_append_error_msg(nullptr, "Not object item", "", nullptr));
898
0
                    ADVANCE_ROW();
899
0
                    continue;
900
0
                }
901
0
                cur = val.get_object();
902
0
            }
903
0
            RETURN_IF_ERROR(_simdjson_write_columns_by_jsonpath(&cur, slot_descs, block, &valid));
904
0
            ADVANCE_ROW();
905
0
            if (!valid) {
906
0
                continue; // process next line
907
0
            }
908
0
            *is_empty_row = false;
909
0
        }
910
0
    } catch (simdjson::simdjson_error& e) {
911
0
        RETURN_IF_ERROR(_handle_simdjson_error(e, block, num_rows, eof));
912
0
        if (*_scanner_eof) {
913
            // When _scanner_eof is true and valid is false, it means that we have encountered
914
            // unqualified data and decided to stop the scan.
915
0
            *is_empty_row = true;
916
0
            return Status::OK();
917
0
        }
918
0
    }
919
920
0
    return Status::OK();
921
0
}
922
923
Status NewJsonReader::_simdjson_handle_nested_complex_json(
924
        RuntimeState* /*state*/, Block& block, const std::vector<SlotDescriptor*>& slot_descs,
925
0
        bool* is_empty_row, bool* eof) {
926
    // nested complex json
927
0
    while (true) {
928
0
        size_t num_rows = block.rows();
929
0
        simdjson::ondemand::object cur;
930
0
        size_t size = 0;
931
0
        simdjson::error_code error;
932
0
        try {
933
0
            RETURN_IF_ERROR(_simdjson_parse_json(&size, is_empty_row, eof, &error));
934
0
            if (size == 0 || *eof) {
935
0
                *is_empty_row = true;
936
0
                return Status::OK();
937
0
            }
938
0
            Status st = _get_json_value(&size, eof, &error, is_empty_row);
939
0
            if (st.is<DATA_QUALITY_ERROR>()) {
940
0
                continue; // continue to read next
941
0
            }
942
0
            RETURN_IF_ERROR(st);
943
0
            if (*is_empty_row) {
944
0
                return Status::OK();
945
0
            }
946
0
            *is_empty_row = false;
947
0
            bool valid = true;
948
0
            if (_json_value.type() != simdjson::ondemand::json_type::object) {
949
0
                RETURN_IF_ERROR(_append_error_msg(nullptr, "Not object item", "", nullptr));
950
0
                continue;
951
0
            }
952
0
            cur = _json_value.get_object();
953
0
            st = _simdjson_write_columns_by_jsonpath(&cur, slot_descs, block, &valid);
954
0
            if (!st.ok()) {
955
0
                RETURN_IF_ERROR(_append_error_msg(nullptr, st.to_string(), "", nullptr));
956
                // Before continuing to process other rows, we need to first clean the fail parsed row.
957
0
                json_reader_detail::truncate_block_to_rows(block, num_rows);
958
0
                continue;
959
0
            }
960
0
            if (!valid) {
961
                // there is only one line in this case, so if it return false, just set is_empty_row true
962
                // so that the caller will continue reading next line.
963
0
                *is_empty_row = true;
964
0
            }
965
0
            break; // read a valid row
966
0
        } catch (simdjson::simdjson_error& e) {
967
0
            RETURN_IF_ERROR(_handle_simdjson_error(e, block, num_rows, eof));
968
0
            if (*_scanner_eof) {
969
                // When _scanner_eof is true and valid is false, it means that we have encountered
970
                // unqualified data and decided to stop the scan.
971
0
                *is_empty_row = true;
972
0
                return Status::OK();
973
0
            }
974
0
            continue;
975
0
        }
976
0
    }
977
0
    return Status::OK();
978
0
}
979
980
2.09k
size_t NewJsonReader::_column_index(const StringRef& name, size_t key_index) {
981
    /// Optimization by caching the order of fields (which is almost always the same)
982
    /// and a quick check to match the next expected field, instead of searching the hash table.
983
2.09k
    if (_prev_positions.size() > key_index && name == _prev_positions[key_index]->first) {
984
0
        return _prev_positions[key_index]->second;
985
0
    }
986
2.09k
    auto it = _slot_desc_index.find(name);
987
2.09k
    if (it != _slot_desc_index.end()) {
988
2.09k
        if (key_index < _prev_positions.size()) {
989
0
            _prev_positions[key_index] = it;
990
0
        }
991
2.09k
        return it->second;
992
2.09k
    }
993
0
    return size_t(-1);
994
2.09k
}
995
996
Status NewJsonReader::_simdjson_set_column_value(simdjson::ondemand::object* value, Block& block,
997
                                                 const std::vector<SlotDescriptor*>& slot_descs,
998
2.09k
                                                 bool* valid) {
999
    // set
1000
2.09k
    _seen_columns.assign(block.columns(), false);
1001
2.09k
    size_t cur_row_count = block.rows();
1002
2.09k
    bool has_valid_value = false;
1003
    // iterate through object, simdjson::ondemond will parsing on the fly
1004
2.09k
    size_t key_index = 0;
1005
1006
2.09k
    for (auto field : *value) {
1007
2.09k
        std::string_view key = field.unescaped_key();
1008
2.09k
        StringRef name_ref(key.data(), key.size());
1009
2.09k
        std::string key_string;
1010
2.09k
        if (_is_hive_table) {
1011
0
            key_string = name_ref.to_string();
1012
0
            std::transform(key_string.begin(), key_string.end(), key_string.begin(), ::tolower);
1013
0
            name_ref = StringRef(key_string);
1014
0
        }
1015
2.09k
        const size_t column_index = _column_index(name_ref, key_index++);
1016
2.09k
        if (UNLIKELY(ssize_t(column_index) < 0)) {
1017
            // This key is not exist in slot desc, just ignore
1018
0
            continue;
1019
0
        }
1020
2.09k
        if (column_index == skip_bitmap_col_idx) {
1021
0
            continue;
1022
0
        }
1023
2.09k
        if (_seen_columns[column_index]) {
1024
0
            if (_is_hive_table) {
1025
                //Since value can only be traversed once,
1026
                // we can only insert the original value first, then delete it, and then reinsert the new value
1027
0
                json_reader_detail::pop_back_last_inserted_value(block, column_index);
1028
0
            } else {
1029
0
                continue;
1030
0
            }
1031
0
        }
1032
2.09k
        simdjson::ondemand::value val = field.value();
1033
2.09k
        auto* column_ptr = block.get_by_position(column_index).column->assert_mutable().get();
1034
2.09k
        RETURN_IF_ERROR(_simdjson_write_data_to_column<false>(
1035
2.09k
                val, slot_descs[column_index]->type(), column_ptr,
1036
2.09k
                slot_descs[column_index]->col_name(), _serdes[column_index], valid));
1037
2.09k
        if (!(*valid)) {
1038
0
            return Status::OK();
1039
0
        }
1040
2.09k
        _seen_columns[column_index] = true;
1041
2.09k
        has_valid_value = true;
1042
2.09k
    }
1043
1044
2.09k
    if (!has_valid_value && _is_load) {
1045
0
        std::string col_names;
1046
0
        for (auto* slot_desc : slot_descs) {
1047
0
            col_names.append(slot_desc->col_name() + ", ");
1048
0
        }
1049
0
        RETURN_IF_ERROR(_append_error_msg(value,
1050
0
                                          "There is no column matching jsonpaths in the json file, "
1051
0
                                          "columns:[{}], please check columns "
1052
0
                                          "and jsonpaths:" +
1053
0
                                                  _jsonpaths,
1054
0
                                          col_names, valid));
1055
0
        return Status::OK();
1056
0
    }
1057
1058
2.09k
    if (_should_process_skip_bitmap_col()) {
1059
0
        _append_empty_skip_bitmap_value(block, cur_row_count);
1060
0
    }
1061
1062
    // fill missing slot
1063
2.09k
    int nullcount = 0;
1064
4.19k
    for (size_t i = 0; i < slot_descs.size(); ++i) {
1065
2.09k
        if (_seen_columns[i]) {
1066
2.09k
            continue;
1067
2.09k
        }
1068
0
        if (i == skip_bitmap_col_idx) {
1069
0
            continue;
1070
0
        }
1071
1072
0
        auto* slot_desc = slot_descs[i];
1073
0
        auto* column_ptr = block.get_by_position(i).column->assert_mutable().get();
1074
1075
        // Quick path to insert default value, instead of using default values in the value map.
1076
0
        if (!_should_process_skip_bitmap_col() &&
1077
0
            (_col_default_value_map.empty() ||
1078
0
             _col_default_value_map.find(slot_desc->col_name()) == _col_default_value_map.end())) {
1079
0
            column_ptr->insert_default();
1080
0
            continue;
1081
0
        }
1082
0
        if (column_ptr->size() < cur_row_count + 1) {
1083
0
            DCHECK(column_ptr->size() == cur_row_count);
1084
0
            if (_should_process_skip_bitmap_col()) {
1085
                // not found, skip this column in flexible partial update
1086
0
                if (slot_desc->is_key() && !slot_desc->is_auto_increment()) {
1087
0
                    RETURN_IF_ERROR(
1088
0
                            _append_error_msg(value,
1089
0
                                              "The key columns can not be ommited in flexible "
1090
0
                                              "partial update, missing key column: {}",
1091
0
                                              slot_desc->col_name(), valid));
1092
                    // remove this line in block
1093
0
                    json_reader_detail::truncate_block_to_rows(block, cur_row_count);
1094
0
                    return Status::OK();
1095
0
                }
1096
0
                _set_skip_bitmap_mark(slot_desc, column_ptr, block, cur_row_count, valid);
1097
0
                column_ptr->insert_default();
1098
0
            } else {
1099
0
                RETURN_IF_ERROR(_fill_missing_column(slot_desc, _serdes[i], column_ptr, valid));
1100
0
                if (!(*valid)) {
1101
0
                    return Status::OK();
1102
0
                }
1103
0
            }
1104
0
            ++nullcount;
1105
0
        }
1106
0
        DCHECK(column_ptr->size() == cur_row_count + 1);
1107
0
    }
1108
1109
    // There is at least one valid value here
1110
2.09k
    DCHECK(nullcount < block.columns());
1111
2.09k
    *valid = true;
1112
2.09k
    return Status::OK();
1113
2.09k
}
1114
1115
template <bool use_string_cache>
1116
Status NewJsonReader::_simdjson_write_data_to_column(simdjson::ondemand::value& value,
1117
                                                     const DataTypePtr& type_desc,
1118
                                                     IColumn* column_ptr,
1119
                                                     const std::string& column_name,
1120
2.09k
                                                     DataTypeSerDeSPtr serde, bool* valid) {
1121
2.09k
    ColumnNullable* nullable_column = nullptr;
1122
2.09k
    IColumn* data_column_ptr = column_ptr;
1123
2.09k
    DataTypeSerDeSPtr data_serde = serde;
1124
1125
2.09k
    if (is_column_nullable(*column_ptr)) {
1126
2.09k
        nullable_column = reinterpret_cast<ColumnNullable*>(column_ptr);
1127
1128
2.09k
        data_column_ptr = nullable_column->get_nested_column().get_ptr().get();
1129
2.09k
        data_serde = serde->get_nested_serdes()[0];
1130
1131
        // kNullType will put 1 into the Null map, so there is no need to push 0 for kNullType.
1132
2.09k
        if (value.type() == simdjson::ondemand::json_type::null) {
1133
0
            nullable_column->insert_default();
1134
0
            *valid = true;
1135
0
            return Status::OK();
1136
0
        }
1137
2.09k
    } else if (value.type() == simdjson::ondemand::json_type::null) [[unlikely]] {
1138
0
        if (_is_load) {
1139
0
            RETURN_IF_ERROR(_append_error_msg(
1140
0
                    nullptr, "Json value is null, but the column `{}` is not nullable.",
1141
0
                    column_name, valid));
1142
0
            return Status::OK();
1143
0
        } else {
1144
0
            return Status::DataQualityError(
1145
0
                    "Json value is null, but the column `{}` is not nullable.", column_name);
1146
0
        }
1147
0
    }
1148
1149
2.09k
    auto primitive_type = type_desc->get_primitive_type();
1150
2.09k
    if (_is_load || !is_complex_type(primitive_type)) {
1151
2.09k
        if (value.type() == simdjson::ondemand::json_type::string) {
1152
0
            std::string_view value_string;
1153
0
            if constexpr (use_string_cache) {
1154
0
                const auto cache_key = value.raw_json().value();
1155
0
                if (_cached_string_values.contains(cache_key)) {
1156
0
                    value_string = _cached_string_values[cache_key];
1157
0
                } else {
1158
0
                    value_string = value.get_string();
1159
0
                    _cached_string_values.emplace(cache_key, value_string);
1160
0
                }
1161
0
            } else {
1162
0
                DCHECK(_cached_string_values.empty());
1163
0
                value_string = value.get_string();
1164
0
            }
1165
1166
0
            Slice slice {value_string.data(), value_string.size()};
1167
0
            RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1168
0
                                                                       _serde_options));
1169
1170
2.09k
        } else if (value.type() == simdjson::ondemand::json_type::boolean) {
1171
0
            const char* str_value = nullptr;
1172
            // insert "1"/"0" , not "true"/"false".
1173
0
            if (value.get_bool()) {
1174
0
                str_value = (char*)"1";
1175
0
            } else {
1176
0
                str_value = (char*)"0";
1177
0
            }
1178
0
            Slice slice {str_value, 1};
1179
0
            RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1180
0
                                                                       _serde_options));
1181
2.09k
        } else {
1182
            // Maybe we can `switch (value->GetType()) case: kNumberType`.
1183
            // Note that `if (value->IsInt())`, but column is FloatColumn.
1184
2.09k
            std::string_view json_str = simdjson::to_json_string(value);
1185
2.09k
            Slice slice {json_str.data(), json_str.size()};
1186
2.09k
            RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1187
2.09k
                                                                       _serde_options));
1188
2.09k
        }
1189
2.09k
    } else if (primitive_type == TYPE_STRUCT) {
1190
0
        if (value.type() != simdjson::ondemand::json_type::object) [[unlikely]] {
1191
0
            return Status::DataQualityError(
1192
0
                    "Json value isn't object, but the column `{}` is struct.", column_name);
1193
0
        }
1194
1195
0
        const auto* type_struct =
1196
0
                assert_cast<const DataTypeStruct*>(remove_nullable(type_desc).get());
1197
0
        auto sub_col_size = type_struct->get_elements().size();
1198
0
        simdjson::ondemand::object struct_value = value.get_object();
1199
0
        auto sub_serdes = data_serde->get_nested_serdes();
1200
0
        auto* struct_column_ptr = assert_cast<ColumnStruct*>(data_column_ptr);
1201
1202
0
        std::map<std::string, size_t> sub_col_name_to_idx;
1203
0
        for (size_t sub_col_idx = 0; sub_col_idx < sub_col_size; sub_col_idx++) {
1204
0
            sub_col_name_to_idx.emplace(type_struct->get_element_name(sub_col_idx), sub_col_idx);
1205
0
        }
1206
0
        std::vector<bool> has_value(sub_col_size, false);
1207
0
        for (simdjson::ondemand::field sub : struct_value) {
1208
0
            std::string_view sub_key_view = sub.unescaped_key();
1209
0
            std::string sub_key(sub_key_view.data(), sub_key_view.length());
1210
0
            std::transform(sub_key.begin(), sub_key.end(), sub_key.begin(), ::tolower);
1211
1212
0
            if (sub_col_name_to_idx.find(sub_key) == sub_col_name_to_idx.end()) [[unlikely]] {
1213
0
                continue;
1214
0
            }
1215
0
            size_t sub_column_idx = sub_col_name_to_idx[sub_key];
1216
0
            auto sub_column_ptr = struct_column_ptr->get_column(sub_column_idx).get_ptr();
1217
1218
0
            if (has_value[sub_column_idx]) [[unlikely]] {
1219
                // Since struct_value can only be traversed once, we can only insert
1220
                // the original value first, then delete it, and then reinsert the new value.
1221
0
                sub_column_ptr->pop_back(1);
1222
0
            }
1223
0
            has_value[sub_column_idx] = true;
1224
1225
0
            const auto& sub_col_type = type_struct->get_element(sub_column_idx);
1226
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<use_string_cache>(
1227
0
                    sub.value(), sub_col_type, sub_column_ptr.get(), column_name + "." + sub_key,
1228
0
                    sub_serdes[sub_column_idx], valid));
1229
0
        }
1230
1231
        // fill missing subcolumn
1232
0
        for (size_t sub_col_idx = 0; sub_col_idx < sub_col_size; sub_col_idx++) {
1233
0
            if (has_value[sub_col_idx]) {
1234
0
                continue;
1235
0
            }
1236
1237
0
            auto sub_column_ptr = struct_column_ptr->get_column(sub_col_idx).get_ptr();
1238
0
            if (is_column_nullable(*sub_column_ptr)) {
1239
0
                sub_column_ptr->insert_default();
1240
0
                continue;
1241
0
            } else [[unlikely]] {
1242
0
                return Status::DataQualityError(
1243
0
                        "Json file structColumn miss field {} and this column isn't nullable.",
1244
0
                        column_name + "." + type_struct->get_element_name(sub_col_idx));
1245
0
            }
1246
0
        }
1247
0
    } else if (primitive_type == TYPE_MAP) {
1248
0
        if (value.type() != simdjson::ondemand::json_type::object) [[unlikely]] {
1249
0
            return Status::DataQualityError("Json value isn't object, but the column `{}` is map.",
1250
0
                                            column_name);
1251
0
        }
1252
0
        simdjson::ondemand::object object_value = value.get_object();
1253
1254
0
        auto sub_serdes = data_serde->get_nested_serdes();
1255
0
        auto* map_column_ptr = assert_cast<ColumnMap*>(data_column_ptr);
1256
1257
0
        size_t field_count = 0;
1258
0
        for (simdjson::ondemand::field member_value : object_value) {
1259
0
            auto f = [](std::string_view key_view, const DataTypePtr& type_desc,
1260
0
                        IColumn* column_ptr, DataTypeSerDeSPtr serde,
1261
0
                        DataTypeSerDe::FormatOptions serde_options, bool* valid) {
1262
0
                auto* data_column_ptr = column_ptr;
1263
0
                auto data_serde = serde;
1264
0
                if (is_column_nullable(*column_ptr)) {
1265
0
                    auto* nullable_column = static_cast<ColumnNullable*>(column_ptr);
1266
1267
0
                    nullable_column->get_null_map_data().push_back(0);
1268
0
                    data_column_ptr = nullable_column->get_nested_column().get_ptr().get();
1269
0
                    data_serde = serde->get_nested_serdes()[0];
1270
0
                }
1271
0
                Slice slice(key_view.data(), key_view.length());
1272
1273
0
                RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1274
0
                                                                           serde_options));
1275
0
                return Status::OK();
1276
0
            };
Unexecuted instantiation: _ZZN5doris13NewJsonReader30_simdjson_write_data_to_columnILb0EEENS_6StatusERN8simdjson8fallback8ondemand5valueERKSt10shared_ptrIKNS_9IDataTypeEEPNS_7IColumnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEES8_INS_13DataTypeSerDeEEPbENKUlSt17basic_string_viewIcSJ_ESD_SF_SP_NSO_13FormatOptionsESQ_E_clESS_SD_SF_SP_ST_SQ_
Unexecuted instantiation: _ZZN5doris13NewJsonReader30_simdjson_write_data_to_columnILb1EEENS_6StatusERN8simdjson8fallback8ondemand5valueERKSt10shared_ptrIKNS_9IDataTypeEEPNS_7IColumnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEES8_INS_13DataTypeSerDeEEPbENKUlSt17basic_string_viewIcSJ_ESD_SF_SP_NSO_13FormatOptionsESQ_E_clESS_SD_SF_SP_ST_SQ_
1277
1278
0
            RETURN_IF_ERROR(f(member_value.unescaped_key(),
1279
0
                              assert_cast<const DataTypeMap*>(remove_nullable(type_desc).get())
1280
0
                                      ->get_key_type(),
1281
0
                              map_column_ptr->get_keys_ptr()->assert_mutable()->get_ptr().get(),
1282
0
                              sub_serdes[0], _serde_options, valid));
1283
1284
0
            simdjson::ondemand::value field_value = member_value.value();
1285
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<use_string_cache>(
1286
0
                    field_value,
1287
0
                    assert_cast<const DataTypeMap*>(remove_nullable(type_desc).get())
1288
0
                            ->get_value_type(),
1289
0
                    map_column_ptr->get_values_ptr()->assert_mutable()->get_ptr().get(),
1290
0
                    column_name + ".value", sub_serdes[1], valid));
1291
0
            field_count++;
1292
0
        }
1293
1294
0
        auto& offsets = map_column_ptr->get_offsets();
1295
0
        offsets.emplace_back(offsets.back() + field_count);
1296
1297
0
    } else if (primitive_type == TYPE_ARRAY) {
1298
0
        if (value.type() != simdjson::ondemand::json_type::array) [[unlikely]] {
1299
0
            return Status::DataQualityError("Json value isn't array, but the column `{}` is array.",
1300
0
                                            column_name);
1301
0
        }
1302
1303
0
        simdjson::ondemand::array array_value = value.get_array();
1304
1305
0
        auto sub_serdes = data_serde->get_nested_serdes();
1306
0
        auto* array_column_ptr = assert_cast<ColumnArray*>(data_column_ptr);
1307
1308
0
        int field_count = 0;
1309
0
        for (simdjson::ondemand::value sub_value : array_value) {
1310
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<use_string_cache>(
1311
0
                    sub_value,
1312
0
                    assert_cast<const DataTypeArray*>(remove_nullable(type_desc).get())
1313
0
                            ->get_nested_type(),
1314
0
                    array_column_ptr->get_data().get_ptr().get(), column_name + ".element",
1315
0
                    sub_serdes[0], valid));
1316
0
            field_count++;
1317
0
        }
1318
0
        auto& offsets = array_column_ptr->get_offsets();
1319
0
        offsets.emplace_back(offsets.back() + field_count);
1320
1321
0
    } else {
1322
0
        return Status::InternalError("Not support load to complex column.");
1323
0
    }
1324
    //We need to finally set the nullmap of column_nullable to keep the size consistent with data_column
1325
2.09k
    if (nullable_column && value.type() != simdjson::ondemand::json_type::null) {
1326
2.09k
        nullable_column->get_null_map_data().push_back(0);
1327
2.09k
    }
1328
2.09k
    *valid = true;
1329
2.09k
    return Status::OK();
1330
2.09k
}
_ZN5doris13NewJsonReader30_simdjson_write_data_to_columnILb0EEENS_6StatusERN8simdjson8fallback8ondemand5valueERKSt10shared_ptrIKNS_9IDataTypeEEPNS_7IColumnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEES8_INS_13DataTypeSerDeEEPb
Line
Count
Source
1120
2.09k
                                                     DataTypeSerDeSPtr serde, bool* valid) {
1121
2.09k
    ColumnNullable* nullable_column = nullptr;
1122
2.09k
    IColumn* data_column_ptr = column_ptr;
1123
2.09k
    DataTypeSerDeSPtr data_serde = serde;
1124
1125
2.09k
    if (is_column_nullable(*column_ptr)) {
1126
2.09k
        nullable_column = reinterpret_cast<ColumnNullable*>(column_ptr);
1127
1128
2.09k
        data_column_ptr = nullable_column->get_nested_column().get_ptr().get();
1129
2.09k
        data_serde = serde->get_nested_serdes()[0];
1130
1131
        // kNullType will put 1 into the Null map, so there is no need to push 0 for kNullType.
1132
2.09k
        if (value.type() == simdjson::ondemand::json_type::null) {
1133
0
            nullable_column->insert_default();
1134
0
            *valid = true;
1135
0
            return Status::OK();
1136
0
        }
1137
2.09k
    } else if (value.type() == simdjson::ondemand::json_type::null) [[unlikely]] {
1138
0
        if (_is_load) {
1139
0
            RETURN_IF_ERROR(_append_error_msg(
1140
0
                    nullptr, "Json value is null, but the column `{}` is not nullable.",
1141
0
                    column_name, valid));
1142
0
            return Status::OK();
1143
0
        } else {
1144
0
            return Status::DataQualityError(
1145
0
                    "Json value is null, but the column `{}` is not nullable.", column_name);
1146
0
        }
1147
0
    }
1148
1149
2.09k
    auto primitive_type = type_desc->get_primitive_type();
1150
2.09k
    if (_is_load || !is_complex_type(primitive_type)) {
1151
2.09k
        if (value.type() == simdjson::ondemand::json_type::string) {
1152
0
            std::string_view value_string;
1153
            if constexpr (use_string_cache) {
1154
                const auto cache_key = value.raw_json().value();
1155
                if (_cached_string_values.contains(cache_key)) {
1156
                    value_string = _cached_string_values[cache_key];
1157
                } else {
1158
                    value_string = value.get_string();
1159
                    _cached_string_values.emplace(cache_key, value_string);
1160
                }
1161
0
            } else {
1162
0
                DCHECK(_cached_string_values.empty());
1163
0
                value_string = value.get_string();
1164
0
            }
1165
1166
0
            Slice slice {value_string.data(), value_string.size()};
1167
0
            RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1168
0
                                                                       _serde_options));
1169
1170
2.09k
        } else if (value.type() == simdjson::ondemand::json_type::boolean) {
1171
0
            const char* str_value = nullptr;
1172
            // insert "1"/"0" , not "true"/"false".
1173
0
            if (value.get_bool()) {
1174
0
                str_value = (char*)"1";
1175
0
            } else {
1176
0
                str_value = (char*)"0";
1177
0
            }
1178
0
            Slice slice {str_value, 1};
1179
0
            RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1180
0
                                                                       _serde_options));
1181
2.09k
        } else {
1182
            // Maybe we can `switch (value->GetType()) case: kNumberType`.
1183
            // Note that `if (value->IsInt())`, but column is FloatColumn.
1184
2.09k
            std::string_view json_str = simdjson::to_json_string(value);
1185
2.09k
            Slice slice {json_str.data(), json_str.size()};
1186
2.09k
            RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1187
2.09k
                                                                       _serde_options));
1188
2.09k
        }
1189
2.09k
    } else if (primitive_type == TYPE_STRUCT) {
1190
0
        if (value.type() != simdjson::ondemand::json_type::object) [[unlikely]] {
1191
0
            return Status::DataQualityError(
1192
0
                    "Json value isn't object, but the column `{}` is struct.", column_name);
1193
0
        }
1194
1195
0
        const auto* type_struct =
1196
0
                assert_cast<const DataTypeStruct*>(remove_nullable(type_desc).get());
1197
0
        auto sub_col_size = type_struct->get_elements().size();
1198
0
        simdjson::ondemand::object struct_value = value.get_object();
1199
0
        auto sub_serdes = data_serde->get_nested_serdes();
1200
0
        auto* struct_column_ptr = assert_cast<ColumnStruct*>(data_column_ptr);
1201
1202
0
        std::map<std::string, size_t> sub_col_name_to_idx;
1203
0
        for (size_t sub_col_idx = 0; sub_col_idx < sub_col_size; sub_col_idx++) {
1204
0
            sub_col_name_to_idx.emplace(type_struct->get_element_name(sub_col_idx), sub_col_idx);
1205
0
        }
1206
0
        std::vector<bool> has_value(sub_col_size, false);
1207
0
        for (simdjson::ondemand::field sub : struct_value) {
1208
0
            std::string_view sub_key_view = sub.unescaped_key();
1209
0
            std::string sub_key(sub_key_view.data(), sub_key_view.length());
1210
0
            std::transform(sub_key.begin(), sub_key.end(), sub_key.begin(), ::tolower);
1211
1212
0
            if (sub_col_name_to_idx.find(sub_key) == sub_col_name_to_idx.end()) [[unlikely]] {
1213
0
                continue;
1214
0
            }
1215
0
            size_t sub_column_idx = sub_col_name_to_idx[sub_key];
1216
0
            auto sub_column_ptr = struct_column_ptr->get_column(sub_column_idx).get_ptr();
1217
1218
0
            if (has_value[sub_column_idx]) [[unlikely]] {
1219
                // Since struct_value can only be traversed once, we can only insert
1220
                // the original value first, then delete it, and then reinsert the new value.
1221
0
                sub_column_ptr->pop_back(1);
1222
0
            }
1223
0
            has_value[sub_column_idx] = true;
1224
1225
0
            const auto& sub_col_type = type_struct->get_element(sub_column_idx);
1226
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<use_string_cache>(
1227
0
                    sub.value(), sub_col_type, sub_column_ptr.get(), column_name + "." + sub_key,
1228
0
                    sub_serdes[sub_column_idx], valid));
1229
0
        }
1230
1231
        // fill missing subcolumn
1232
0
        for (size_t sub_col_idx = 0; sub_col_idx < sub_col_size; sub_col_idx++) {
1233
0
            if (has_value[sub_col_idx]) {
1234
0
                continue;
1235
0
            }
1236
1237
0
            auto sub_column_ptr = struct_column_ptr->get_column(sub_col_idx).get_ptr();
1238
0
            if (is_column_nullable(*sub_column_ptr)) {
1239
0
                sub_column_ptr->insert_default();
1240
0
                continue;
1241
0
            } else [[unlikely]] {
1242
0
                return Status::DataQualityError(
1243
0
                        "Json file structColumn miss field {} and this column isn't nullable.",
1244
0
                        column_name + "." + type_struct->get_element_name(sub_col_idx));
1245
0
            }
1246
0
        }
1247
0
    } else if (primitive_type == TYPE_MAP) {
1248
0
        if (value.type() != simdjson::ondemand::json_type::object) [[unlikely]] {
1249
0
            return Status::DataQualityError("Json value isn't object, but the column `{}` is map.",
1250
0
                                            column_name);
1251
0
        }
1252
0
        simdjson::ondemand::object object_value = value.get_object();
1253
1254
0
        auto sub_serdes = data_serde->get_nested_serdes();
1255
0
        auto* map_column_ptr = assert_cast<ColumnMap*>(data_column_ptr);
1256
1257
0
        size_t field_count = 0;
1258
0
        for (simdjson::ondemand::field member_value : object_value) {
1259
0
            auto f = [](std::string_view key_view, const DataTypePtr& type_desc,
1260
0
                        IColumn* column_ptr, DataTypeSerDeSPtr serde,
1261
0
                        DataTypeSerDe::FormatOptions serde_options, bool* valid) {
1262
0
                auto* data_column_ptr = column_ptr;
1263
0
                auto data_serde = serde;
1264
0
                if (is_column_nullable(*column_ptr)) {
1265
0
                    auto* nullable_column = static_cast<ColumnNullable*>(column_ptr);
1266
1267
0
                    nullable_column->get_null_map_data().push_back(0);
1268
0
                    data_column_ptr = nullable_column->get_nested_column().get_ptr().get();
1269
0
                    data_serde = serde->get_nested_serdes()[0];
1270
0
                }
1271
0
                Slice slice(key_view.data(), key_view.length());
1272
1273
0
                RETURN_IF_ERROR(data_serde->deserialize_one_cell_from_json(*data_column_ptr, slice,
1274
0
                                                                           serde_options));
1275
0
                return Status::OK();
1276
0
            };
1277
1278
0
            RETURN_IF_ERROR(f(member_value.unescaped_key(),
1279
0
                              assert_cast<const DataTypeMap*>(remove_nullable(type_desc).get())
1280
0
                                      ->get_key_type(),
1281
0
                              map_column_ptr->get_keys_ptr()->assert_mutable()->get_ptr().get(),
1282
0
                              sub_serdes[0], _serde_options, valid));
1283
1284
0
            simdjson::ondemand::value field_value = member_value.value();
1285
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<use_string_cache>(
1286
0
                    field_value,
1287
0
                    assert_cast<const DataTypeMap*>(remove_nullable(type_desc).get())
1288
0
                            ->get_value_type(),
1289
0
                    map_column_ptr->get_values_ptr()->assert_mutable()->get_ptr().get(),
1290
0
                    column_name + ".value", sub_serdes[1], valid));
1291
0
            field_count++;
1292
0
        }
1293
1294
0
        auto& offsets = map_column_ptr->get_offsets();
1295
0
        offsets.emplace_back(offsets.back() + field_count);
1296
1297
0
    } else if (primitive_type == TYPE_ARRAY) {
1298
0
        if (value.type() != simdjson::ondemand::json_type::array) [[unlikely]] {
1299
0
            return Status::DataQualityError("Json value isn't array, but the column `{}` is array.",
1300
0
                                            column_name);
1301
0
        }
1302
1303
0
        simdjson::ondemand::array array_value = value.get_array();
1304
1305
0
        auto sub_serdes = data_serde->get_nested_serdes();
1306
0
        auto* array_column_ptr = assert_cast<ColumnArray*>(data_column_ptr);
1307
1308
0
        int field_count = 0;
1309
0
        for (simdjson::ondemand::value sub_value : array_value) {
1310
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<use_string_cache>(
1311
0
                    sub_value,
1312
0
                    assert_cast<const DataTypeArray*>(remove_nullable(type_desc).get())
1313
0
                            ->get_nested_type(),
1314
0
                    array_column_ptr->get_data().get_ptr().get(), column_name + ".element",
1315
0
                    sub_serdes[0], valid));
1316
0
            field_count++;
1317
0
        }
1318
0
        auto& offsets = array_column_ptr->get_offsets();
1319
0
        offsets.emplace_back(offsets.back() + field_count);
1320
1321
0
    } else {
1322
0
        return Status::InternalError("Not support load to complex column.");
1323
0
    }
1324
    //We need to finally set the nullmap of column_nullable to keep the size consistent with data_column
1325
2.09k
    if (nullable_column && value.type() != simdjson::ondemand::json_type::null) {
1326
2.09k
        nullable_column->get_null_map_data().push_back(0);
1327
2.09k
    }
1328
2.09k
    *valid = true;
1329
2.09k
    return Status::OK();
1330
2.09k
}
Unexecuted instantiation: _ZN5doris13NewJsonReader30_simdjson_write_data_to_columnILb1EEENS_6StatusERN8simdjson8fallback8ondemand5valueERKSt10shared_ptrIKNS_9IDataTypeEEPNS_7IColumnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEES8_INS_13DataTypeSerDeEEPb
1331
1332
Status NewJsonReader::_append_error_msg(simdjson::ondemand::object* obj, std::string error_msg,
1333
0
                                        std::string col_name, bool* valid) {
1334
0
    std::string err_msg;
1335
0
    if (!col_name.empty()) {
1336
0
        fmt::memory_buffer error_buf;
1337
0
        fmt::format_to(error_buf, error_msg, col_name, _jsonpaths);
1338
0
        err_msg = fmt::to_string(error_buf);
1339
0
    } else {
1340
0
        err_msg = error_msg;
1341
0
    }
1342
1343
0
    _counter->num_rows_filtered++;
1344
0
    if (valid != nullptr) {
1345
        // current row is invalid
1346
0
        *valid = false;
1347
0
    }
1348
1349
0
    RETURN_IF_ERROR(_state->append_error_msg_to_file(
1350
0
            [&]() -> std::string {
1351
0
                if (!obj) {
1352
0
                    return "";
1353
0
                }
1354
0
                std::string_view str_view;
1355
0
                (void)!obj->raw_json().get(str_view);
1356
0
                return std::string(str_view.data(), str_view.size());
1357
0
            },
1358
0
            [&]() -> std::string { return err_msg; }));
1359
0
    return Status::OK();
1360
0
}
1361
1362
Status NewJsonReader::_simdjson_parse_json(size_t* size, bool* is_empty_row, bool* eof,
1363
36.1k
                                           simdjson::error_code* error) {
1364
36.1k
    SCOPED_TIMER(_read_timer);
1365
    // step1: read buf from pipe.
1366
36.1k
    if (_line_reader != nullptr) {
1367
36.1k
        RETURN_IF_ERROR(_line_reader->read_line(&_json_str, size, eof, _io_ctx));
1368
36.1k
    } else {
1369
0
        size_t length = 0;
1370
0
        RETURN_IF_ERROR(_read_one_message(&_json_str_ptr, &length));
1371
0
        _json_str = _json_str_ptr.get();
1372
0
        *size = length;
1373
0
        if (length == 0) {
1374
0
            *eof = true;
1375
0
        }
1376
0
    }
1377
36.1k
    if (*eof) {
1378
1.30k
        return Status::OK();
1379
1.30k
    }
1380
1381
    // step2: init json parser iterate.
1382
34.8k
    if (*size + simdjson::SIMDJSON_PADDING > _padded_size) {
1383
        // For efficiency reasons, simdjson requires a string with a few bytes (simdjson::SIMDJSON_PADDING) at the end.
1384
        // Hence, a re-allocation is needed if the space is not enough.
1385
0
        _simdjson_ondemand_padding_buffer.resize(*size + simdjson::SIMDJSON_PADDING);
1386
0
        _simdjson_ondemand_unscape_padding_buffer.resize(*size + simdjson::SIMDJSON_PADDING);
1387
0
        _padded_size = *size + simdjson::SIMDJSON_PADDING;
1388
0
    }
1389
    // trim BOM since simdjson does not handle UTF-8 Unicode (with BOM)
1390
34.8k
    if (*size >= 3 && static_cast<char>(_json_str[0]) == '\xEF' &&
1391
34.8k
        static_cast<char>(_json_str[1]) == '\xBB' && static_cast<char>(_json_str[2]) == '\xBF') {
1392
        // skip the first three BOM bytes
1393
0
        _json_str += 3;
1394
0
        *size -= 3;
1395
0
    }
1396
34.8k
    memcpy(&_simdjson_ondemand_padding_buffer.front(), _json_str, *size);
1397
34.8k
    _original_doc_size = *size;
1398
34.8k
    *error = _ondemand_json_parser
1399
34.8k
                     ->iterate(std::string_view(_simdjson_ondemand_padding_buffer.data(), *size),
1400
34.8k
                               _padded_size)
1401
34.8k
                     .get(_original_json_doc);
1402
34.8k
    return Status::OK();
1403
36.1k
}
1404
1405
2.09k
Status NewJsonReader::_judge_empty_row(size_t size, bool eof, bool* is_empty_row) {
1406
2.09k
    if (size == 0 || eof) {
1407
0
        *is_empty_row = true;
1408
0
        return Status::OK();
1409
0
    }
1410
1411
2.09k
    if (!_parsed_jsonpaths.empty() && _strip_outer_array) {
1412
0
        _total_rows = _json_value.count_elements().value();
1413
0
        _next_row = 0;
1414
1415
0
        if (_total_rows == 0) {
1416
            // meet an empty json array.
1417
0
            *is_empty_row = true;
1418
0
        }
1419
0
    }
1420
2.09k
    return Status::OK();
1421
2.09k
}
1422
1423
Status NewJsonReader::_get_json_value(size_t* size, bool* eof, simdjson::error_code* error,
1424
2.09k
                                      bool* is_empty_row) {
1425
2.09k
    SCOPED_TIMER(_read_timer);
1426
2.09k
    auto return_quality_error = [&](fmt::memory_buffer& error_msg,
1427
2.09k
                                    const std::string& doc_info) -> Status {
1428
0
        _counter->num_rows_filtered++;
1429
0
        RETURN_IF_ERROR(_state->append_error_msg_to_file(
1430
0
                [&]() -> std::string { return doc_info; },
1431
0
                [&]() -> std::string { return fmt::to_string(error_msg); }));
1432
0
        if (*_scanner_eof) {
1433
            // Case A: if _scanner_eof is set to true in "append_error_msg_to_file", which means
1434
            // we meet enough invalid rows and the scanner should be stopped.
1435
            // So we set eof to true and return OK, the caller will stop the process as we meet the end of file.
1436
0
            *eof = true;
1437
0
            return Status::OK();
1438
0
        }
1439
0
        return Status::DataQualityError(fmt::to_string(error_msg));
1440
0
    };
1441
2.09k
    if (*error != simdjson::error_code::SUCCESS) {
1442
0
        fmt::memory_buffer error_msg;
1443
0
        fmt::format_to(error_msg, "Parse json data for JsonDoc failed. code: {}, error info: {}",
1444
0
                       *error, simdjson::error_message(*error));
1445
0
        return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1446
0
    }
1447
2.09k
    auto type_res = _original_json_doc.type();
1448
2.09k
    if (type_res.error() != simdjson::error_code::SUCCESS) {
1449
0
        fmt::memory_buffer error_msg;
1450
0
        fmt::format_to(error_msg, "Parse json data for JsonDoc failed. code: {}, error info: {}",
1451
0
                       type_res.error(), simdjson::error_message(type_res.error()));
1452
0
        return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1453
0
    }
1454
2.09k
    simdjson::ondemand::json_type type = type_res.value();
1455
2.09k
    if (type != simdjson::ondemand::json_type::object &&
1456
2.09k
        type != simdjson::ondemand::json_type::array) {
1457
0
        fmt::memory_buffer error_msg;
1458
0
        fmt::format_to(error_msg, "Not an json object or json array");
1459
0
        return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1460
0
    }
1461
2.09k
    if (!_parsed_json_root.empty() && type == simdjson::ondemand::json_type::object) {
1462
0
        try {
1463
            // set json root
1464
            // if it is an array at top level, then we should iterate the entire array in
1465
            // ::_simdjson_handle_flat_array_complex_json
1466
0
            simdjson::ondemand::object object = _original_json_doc;
1467
0
            Status st = JsonFunctions::extract_from_object(object, _parsed_json_root, &_json_value);
1468
0
            if (!st.ok()) {
1469
0
                fmt::memory_buffer error_msg;
1470
0
                fmt::format_to(error_msg, "{}", st.to_string());
1471
0
                return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1472
0
            }
1473
0
            _parsed_from_json_root = true;
1474
0
        } catch (simdjson::simdjson_error& e) {
1475
0
            fmt::memory_buffer error_msg;
1476
0
            fmt::format_to(error_msg, "Encounter error while extract_from_object, error: {}",
1477
0
                           e.what());
1478
0
            return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1479
0
        }
1480
2.09k
    } else {
1481
2.09k
        _json_value = _original_json_doc;
1482
2.09k
    }
1483
1484
2.09k
    if (_json_value.type() == simdjson::ondemand::json_type::array && !_strip_outer_array) {
1485
0
        fmt::memory_buffer error_msg;
1486
0
        fmt::format_to(error_msg, "{}",
1487
0
                       "JSON data is array-object, `strip_outer_array` must be TRUE.");
1488
0
        return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1489
0
    }
1490
1491
2.09k
    if (_json_value.type() != simdjson::ondemand::json_type::array && _strip_outer_array) {
1492
0
        fmt::memory_buffer error_msg;
1493
0
        fmt::format_to(error_msg, "{}",
1494
0
                       "JSON data is not an array-object, `strip_outer_array` must be FALSE.");
1495
0
        return return_quality_error(error_msg, std::string((char*)_json_str, *size));
1496
0
    }
1497
2.09k
    RETURN_IF_ERROR(_judge_empty_row(*size, *eof, is_empty_row));
1498
2.09k
    return Status::OK();
1499
2.09k
}
1500
1501
Status NewJsonReader::_simdjson_write_columns_by_jsonpath(
1502
        simdjson::ondemand::object* value, const std::vector<SlotDescriptor*>& slot_descs,
1503
0
        Block& block, bool* valid) {
1504
    // write by jsonpath
1505
0
    bool has_valid_value = false;
1506
1507
0
    Defer clear_defer([this]() { _cached_string_values.clear(); });
1508
1509
0
    for (size_t i = 0; i < slot_descs.size(); i++) {
1510
0
        auto* slot_desc = slot_descs[i];
1511
0
        auto* column_ptr = block.get_by_position(i).column->assert_mutable().get();
1512
0
        simdjson::ondemand::value json_value;
1513
0
        Status st;
1514
0
        if (i < _parsed_jsonpaths.size()) {
1515
0
            st = JsonFunctions::extract_from_object(*value, _parsed_jsonpaths[i], &json_value);
1516
0
            if (!st.ok() && !st.is<NOT_FOUND>()) {
1517
0
                return st;
1518
0
            }
1519
0
        }
1520
0
        if (i < _parsed_jsonpaths.size() && JsonFunctions::is_root_path(_parsed_jsonpaths[i])) {
1521
            // Indicate that the jsonpath is "$" or "$.", read the full root json object, insert the original doc directly
1522
0
            ColumnNullable* nullable_column = nullptr;
1523
0
            IColumn* target_column_ptr = nullptr;
1524
0
            if (slot_desc->is_nullable()) {
1525
0
                nullable_column = assert_cast<ColumnNullable*>(column_ptr);
1526
0
                target_column_ptr = &nullable_column->get_nested_column();
1527
0
                nullable_column->get_null_map_data().push_back(0);
1528
0
            }
1529
0
            auto* column_string = assert_cast<ColumnString*>(target_column_ptr);
1530
0
            column_string->insert_data(_simdjson_ondemand_padding_buffer.data(),
1531
0
                                       _original_doc_size);
1532
0
            has_valid_value = true;
1533
0
        } else if (i >= _parsed_jsonpaths.size() || st.is<NOT_FOUND>()) {
1534
            // not match in jsondata, filling with default value
1535
0
            RETURN_IF_ERROR(_fill_missing_column(slot_desc, _serdes[i], column_ptr, valid));
1536
0
            if (!(*valid)) {
1537
0
                return Status::OK();
1538
0
            }
1539
0
        } else {
1540
0
            RETURN_IF_ERROR(_simdjson_write_data_to_column<true>(json_value, slot_desc->type(),
1541
0
                                                                 column_ptr, slot_desc->col_name(),
1542
0
                                                                 _serdes[i], valid));
1543
0
            if (!(*valid)) {
1544
0
                return Status::OK();
1545
0
            }
1546
0
            has_valid_value = true;
1547
0
        }
1548
0
    }
1549
0
    if (!has_valid_value) {
1550
        // there is no valid value in json line but has filled with default value before
1551
        // so remove this line in block
1552
0
        std::string col_names;
1553
0
        DCHECK(block.rows() > 0);
1554
0
        json_reader_detail::truncate_block_to_rows(block, block.rows() - 1);
1555
0
        for (auto* slot_desc : slot_descs) {
1556
0
            col_names.append(slot_desc->col_name() + ", ");
1557
0
        }
1558
0
        RETURN_IF_ERROR(_append_error_msg(value,
1559
0
                                          "There is no column matching jsonpaths in the json file, "
1560
0
                                          "columns:[{}], please check columns "
1561
0
                                          "and jsonpaths:" +
1562
0
                                                  _jsonpaths,
1563
0
                                          col_names, valid));
1564
0
        return Status::OK();
1565
0
    }
1566
0
    *valid = true;
1567
0
    return Status::OK();
1568
0
}
1569
1570
Status NewJsonReader::_get_column_default_value(
1571
        const std::vector<SlotDescriptor*>& slot_descs,
1572
1.30k
        const std::unordered_map<std::string, VExprContextSPtr>& col_default_value_ctx) {
1573
1.30k
    for (auto* slot_desc : slot_descs) {
1574
1.30k
        auto it = col_default_value_ctx.find(slot_desc->col_name());
1575
1.30k
        if (it != col_default_value_ctx.end() && it->second != nullptr) {
1576
0
            const auto& ctx = it->second;
1577
            // NULL_LITERAL means no valid value of current column
1578
0
            if (ctx->root()->node_type() == TExprNodeType::type::NULL_LITERAL) {
1579
0
                continue;
1580
0
            }
1581
0
            ColumnWithTypeAndName result;
1582
0
            RETURN_IF_ERROR(ctx->execute_const_expr(result));
1583
0
            DCHECK(result.column->size() == 1);
1584
0
            _col_default_value_map.emplace(slot_desc->col_name(),
1585
0
                                           result.column->get_data_at(0).to_string());
1586
0
        }
1587
1.30k
    }
1588
1.30k
    return Status::OK();
1589
1.30k
}
1590
1591
Status NewJsonReader::_fill_missing_column(SlotDescriptor* slot_desc, DataTypeSerDeSPtr serde,
1592
0
                                           IColumn* column_ptr, bool* valid) {
1593
0
    auto col_value = _col_default_value_map.find(slot_desc->col_name());
1594
0
    if (col_value == _col_default_value_map.end()) {
1595
0
        if (slot_desc->is_nullable()) {
1596
0
            auto* nullable_column = static_cast<ColumnNullable*>(column_ptr);
1597
0
            nullable_column->insert_default();
1598
0
        } else {
1599
0
            if (_is_load) {
1600
0
                RETURN_IF_ERROR(_append_error_msg(
1601
0
                        nullptr, "The column `{}` is not nullable, but it's not found in jsondata.",
1602
0
                        slot_desc->col_name(), valid));
1603
0
            } else {
1604
0
                return Status::DataQualityError(
1605
0
                        "The column `{}` is not nullable, but it's not found in jsondata.",
1606
0
                        slot_desc->col_name());
1607
0
            }
1608
0
        }
1609
0
    } else {
1610
0
        const std::string& v_str = col_value->second;
1611
0
        Slice column_default_value {v_str};
1612
0
        RETURN_IF_ERROR(serde->deserialize_one_cell_from_json(*column_ptr, column_default_value,
1613
0
                                                              _serde_options));
1614
0
    }
1615
0
    *valid = true;
1616
0
    return Status::OK();
1617
0
}
1618
1619
0
void NewJsonReader::_append_empty_skip_bitmap_value(Block& block, size_t cur_row_count) {
1620
0
    auto* skip_bitmap_nullable_col_ptr = assert_cast<ColumnNullable*>(
1621
0
            block.get_by_position(skip_bitmap_col_idx).column->assert_mutable().get());
1622
0
    auto* skip_bitmap_col_ptr =
1623
0
            assert_cast<ColumnBitmap*>(skip_bitmap_nullable_col_ptr->get_nested_column_ptr().get());
1624
0
    DCHECK(skip_bitmap_nullable_col_ptr->size() == cur_row_count);
1625
    // should append an empty bitmap for every row wheather this line misses columns
1626
0
    skip_bitmap_nullable_col_ptr->get_null_map_data().push_back(0);
1627
0
    skip_bitmap_col_ptr->insert_default();
1628
0
    DCHECK(skip_bitmap_col_ptr->size() == cur_row_count + 1);
1629
0
}
1630
1631
void NewJsonReader::_set_skip_bitmap_mark(SlotDescriptor* slot_desc, IColumn* column_ptr,
1632
0
                                          Block& block, size_t cur_row_count, bool* valid) {
1633
    // we record the missing column's column unique id in skip bitmap
1634
    // to indicate which columns need to do the alignment process
1635
0
    auto* skip_bitmap_nullable_col_ptr = assert_cast<ColumnNullable*>(
1636
0
            block.get_by_position(skip_bitmap_col_idx).column->assert_mutable().get());
1637
0
    auto* skip_bitmap_col_ptr =
1638
0
            assert_cast<ColumnBitmap*>(skip_bitmap_nullable_col_ptr->get_nested_column_ptr().get());
1639
0
    DCHECK(skip_bitmap_col_ptr->size() == cur_row_count + 1);
1640
0
    auto& skip_bitmap = skip_bitmap_col_ptr->get_data().back();
1641
0
    skip_bitmap.add(slot_desc->col_unique_id());
1642
0
}
1643
1644
0
void NewJsonReader::_collect_profile_before_close() {
1645
0
    if (_line_reader != nullptr) {
1646
0
        _line_reader->collect_profile_before_close();
1647
0
    }
1648
0
    if (_file_reader != nullptr) {
1649
0
        _file_reader->collect_profile_before_close();
1650
0
    }
1651
0
}
1652
1653
} // namespace doris