Coverage Report

Created: 2026-09-28 19:35

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/csv/csv_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/csv/csv_reader.h"
19
20
#include <fmt/format.h>
21
#include <gen_cpp/PlanNodes_types.h>
22
#include <gen_cpp/Types_types.h>
23
#include <glog/logging.h>
24
25
#include <algorithm>
26
#include <cstddef>
27
#include <map>
28
#include <memory>
29
#include <numeric>
30
#include <ostream>
31
#include <regex>
32
#include <utility>
33
34
#include "common/compiler_util.h" // IWYU pragma: keep
35
#include "common/config.h"
36
#include "common/consts.h"
37
#include "common/status.h"
38
#include "core/block/block.h"
39
#include "core/block/column_with_type_and_name.h"
40
#include "core/column/column_nullable.h"
41
#include "core/column/column_string.h"
42
#include "core/data_type/data_type_factory.hpp"
43
#include "core/data_type_serde/data_type_string_serde.h"
44
#include "exec/scan/scanner.h"
45
#include "format/file_reader/new_plain_binary_line_reader.h"
46
#include "format/file_reader/new_plain_text_line_reader.h"
47
#include "format/line_reader.h"
48
#include "io/file_factory.h"
49
#include "io/fs/broker_file_reader.h"
50
#include "io/fs/buffered_reader.h"
51
#include "io/fs/file_reader.h"
52
#include "io/fs/s3_file_reader.h"
53
#include "io/fs/tracing_file_reader.h"
54
#include "runtime/descriptors.h"
55
#include "runtime/runtime_state.h"
56
#include "util/decompressor.h"
57
#include "util/string_util.h"
58
#include "util/utf8_check.h"
59
60
namespace doris {
61
class RuntimeProfile;
62
class IColumn;
63
namespace io {
64
struct IOContext;
65
enum class FileCachePolicy : uint8_t;
66
} // namespace io
67
} // namespace doris
68
69
namespace doris {
70
71
namespace {
72
73
10.8M
size_t columns_byte_size(const std::vector<MutableColumnPtr>& columns) {
74
10.8M
    size_t bytes = 0;
75
166M
    for (const auto& column : columns) {
76
166M
        DCHECK(column.get() != nullptr);
77
166M
        bytes += column->byte_size();
78
166M
    }
79
10.8M
    return bytes;
80
10.8M
}
81
82
} // namespace
83
84
708
void EncloseCsvTextFieldSplitter::do_split(const Slice& line, std::vector<Slice>* splitted_values) {
85
708
    const char* data = line.data;
86
708
    const auto& column_sep_positions = _text_line_reader_ctx->column_sep_positions();
87
708
    size_t value_start_offset = 0;
88
2.57k
    for (auto idx : column_sep_positions) {
89
2.57k
        process_value_func(data, value_start_offset, idx - value_start_offset, _trimming_char,
90
2.57k
                           splitted_values);
91
2.57k
        value_start_offset = idx + _value_sep_len;
92
2.57k
    }
93
708
    if (line.size >= value_start_offset) {
94
        // process the last column
95
708
        process_value_func(data, value_start_offset, line.size - value_start_offset, _trimming_char,
96
708
                           splitted_values);
97
708
    }
98
708
}
99
100
void PlainCsvTextFieldSplitter::_split_field_single_char(const Slice& line,
101
10.8M
                                                         std::vector<Slice>* splitted_values) {
102
10.8M
    const char* data = line.data;
103
10.8M
    const size_t size = line.size;
104
10.8M
    size_t value_start = 0;
105
1.40G
    for (size_t i = 0; i < size; ++i) {
106
1.39G
        if (data[i] == _value_sep[0]) {
107
167M
            process_value_func(data, value_start, i - value_start, _trimming_char, splitted_values);
108
167M
            value_start = i + _value_sep_len;
109
167M
        }
110
1.39G
    }
111
10.8M
    process_value_func(data, value_start, size - value_start, _trimming_char, splitted_values);
112
10.8M
}
113
114
void PlainCsvTextFieldSplitter::_split_field_multi_char(const Slice& line,
115
66
                                                        std::vector<Slice>* splitted_values) {
116
66
    size_t start = 0;  // point to the start pos of next col value.
117
66
    size_t curpos = 0; // point to the start pos of separator matching sequence.
118
119
    // value_sep : AAAA
120
    // line.data : 1234AAAA5678
121
    // -> 1234,5678
122
123
    //    start   start
124
    //      ▼       ▼
125
    //      1234AAAA5678\0
126
    //          ▲       ▲
127
    //      curpos     curpos
128
129
    //kmp
130
66
    std::vector<int> next(_value_sep_len);
131
66
    next[0] = -1;
132
211
    for (int i = 1, j = -1; i < _value_sep_len; i++) {
133
165
        while (j > -1 && _value_sep[i] != _value_sep[j + 1]) {
134
20
            j = next[j];
135
20
        }
136
145
        if (_value_sep[i] == _value_sep[j + 1]) {
137
64
            j++;
138
64
        }
139
145
        next[i] = j;
140
145
    }
141
142
2.86k
    for (int i = 0, j = -1; i < line.size; i++) {
143
        // i : line
144
        // j : _value_sep
145
2.95k
        while (j > -1 && line[i] != _value_sep[j + 1]) {
146
154
            j = next[j];
147
154
        }
148
2.80k
        if (line[i] == _value_sep[j + 1]) {
149
1.02k
            j++;
150
1.02k
        }
151
2.80k
        if (j == _value_sep_len - 1) {
152
270
            curpos = i - _value_sep_len + 1;
153
154
            /*
155
             * column_separator : "xx"
156
             * data.csv :  data1xxxxdata2
157
             *
158
             * Parse incorrectly:
159
             *      data1[xx]xxdata2
160
             *      data1x[xx]xdata2
161
             *      data1xx[xx]data2
162
             * The string "xxxx" is parsed into three "xx" delimiters.
163
             *
164
             * Parse correctly:
165
             *      data1[xx]xxdata2
166
             *      data1xx[xx]data2
167
             */
168
169
270
            if (curpos >= start) {
170
243
                process_value_func(line.data, start, curpos - start, _trimming_char,
171
243
                                   splitted_values);
172
243
                start = i + 1;
173
243
            }
174
175
270
            j = next[j];
176
270
        }
177
2.80k
    }
178
66
    process_value_func(line.data, start, line.size - start, _trimming_char, splitted_values);
179
66
}
180
181
10.8M
void PlainCsvTextFieldSplitter::do_split(const Slice& line, std::vector<Slice>* splitted_values) {
182
10.8M
    if (is_single_char_delim) {
183
10.8M
        _split_field_single_char(line, splitted_values);
184
18.4E
    } else {
185
18.4E
        _split_field_multi_char(line, splitted_values);
186
18.4E
    }
187
10.8M
}
188
189
CsvReader::CsvReader(RuntimeState* state, RuntimeProfile* profile, ScannerCounter* counter,
190
                     const TFileScanRangeParams& params, const TFileRangeDesc& range,
191
                     const std::vector<SlotDescriptor*>& file_slot_descs, size_t batch_size,
192
                     io::IOContext* io_ctx, std::shared_ptr<io::IOContext> io_ctx_holder)
193
3.08k
        : _profile(profile),
194
3.08k
          _params(params),
195
3.08k
          _file_reader(nullptr),
196
3.08k
          _line_reader(nullptr),
197
3.08k
          _decompressor(nullptr),
198
3.08k
          _state(state),
199
3.08k
          _counter(counter),
200
3.08k
          _range(range),
201
3.08k
          _file_slot_descs(file_slot_descs),
202
3.08k
          _line_reader_eof(false),
203
3.08k
          _skip_lines(0),
204
3.08k
          _io_ctx(io_ctx),
205
3.08k
          _io_ctx_holder(std::move(io_ctx_holder)),
206
3.08k
          _batch_size(std::max(batch_size, 1UL)) {
207
3.08k
    if (_io_ctx == nullptr && _io_ctx_holder) {
208
2.08k
        _io_ctx = _io_ctx_holder.get();
209
2.08k
    }
210
3.08k
    _file_format_type = _params.format_type;
211
3.08k
    _is_proto_format = _file_format_type == TFileFormatType::FORMAT_PROTO;
212
3.08k
    if (_range.__isset.compress_type) {
213
        // for compatibility
214
403
        _file_compress_type = _range.compress_type;
215
2.68k
    } else {
216
2.68k
        _file_compress_type = _params.compress_type;
217
2.68k
    }
218
3.08k
    _size = _range.size;
219
220
3.08k
    _split_values.reserve(_file_slot_descs.size());
221
3.08k
    _init_system_properties();
222
3.08k
    _init_file_description();
223
3.08k
    _serdes = create_data_type_serdes(_file_slot_descs);
224
3.08k
}
225
226
3.08k
void CsvReader::_init_system_properties() {
227
3.08k
    if (_range.__isset.file_type) {
228
        // for compatibility
229
403
        _system_properties.system_type = _range.file_type;
230
2.68k
    } else {
231
2.68k
        _system_properties.system_type = _params.file_type;
232
2.68k
    }
233
3.08k
    _system_properties.properties = _params.properties;
234
3.08k
    _system_properties.hdfs_params = _params.hdfs_params;
235
3.08k
    if (_params.__isset.broker_addresses) {
236
2.08k
        _system_properties.broker_addresses.assign(_params.broker_addresses.begin(),
237
2.08k
                                                   _params.broker_addresses.end());
238
2.08k
    }
239
3.08k
}
240
241
3.08k
void CsvReader::_init_file_description() {
242
3.08k
    _file_description.path = _range.path;
243
3.08k
    _file_description.file_size = _range.__isset.file_size ? _range.file_size : -1;
244
3.08k
    if (_range.__isset.fs_name) {
245
0
        _file_description.fs_name = _range.fs_name;
246
0
    }
247
3.08k
    if (_range.__isset.file_cache_admission) {
248
0
        _file_description.file_cache_admission = _range.file_cache_admission;
249
0
    }
250
3.08k
}
251
252
598
Status CsvReader::init_reader(bool is_load) {
253
    // set the skip lines and start offset
254
598
    _start_offset = _range.start_offset;
255
598
    if (_start_offset == 0) {
256
        // check header typer first
257
228
        if (_params.__isset.file_attributes && _params.file_attributes.__isset.header_type &&
258
228
            !_params.file_attributes.header_type.empty()) {
259
4
            std::string header_type = to_lower(_params.file_attributes.header_type);
260
4
            if (header_type == BeConsts::CSV_WITH_NAMES) {
261
2
                _skip_lines = 1;
262
2
            } else if (header_type == BeConsts::CSV_WITH_NAMES_AND_TYPES) {
263
2
                _skip_lines = 2;
264
2
            }
265
224
        } else if (_params.file_attributes.__isset.skip_lines) {
266
2
            _skip_lines = _params.file_attributes.skip_lines;
267
2
        }
268
370
    } else if (_start_offset != 0) {
269
370
        if ((_file_compress_type != TFileCompressType::PLAIN) ||
270
370
            (_file_compress_type == TFileCompressType::UNKNOWN &&
271
370
             _file_format_type != TFileFormatType::FORMAT_CSV_PLAIN)) {
272
0
            return Status::InternalError<false>("For now we do not support split compressed file");
273
0
        }
274
        // pre-read to promise first line skipped always read
275
370
        int64_t pre_read_len = std::min(
276
370
                static_cast<int64_t>(_params.file_attributes.text_params.line_delimiter.size()),
277
370
                _start_offset);
278
370
        _start_offset -= pre_read_len;
279
370
        _size += pre_read_len;
280
        // not first range will always skip one line
281
370
        _skip_lines = 1;
282
370
    }
283
284
598
    _use_nullable_string_opt.resize(_file_slot_descs.size());
285
1.19k
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
286
598
        auto data_type_ptr = _file_slot_descs[i]->get_data_type_ptr();
287
598
        if (data_type_ptr->is_nullable() && is_string_type(data_type_ptr->get_primitive_type())) {
288
598
            _use_nullable_string_opt[i] = 1;
289
598
        }
290
598
    }
291
292
598
    RETURN_IF_ERROR(_init_options());
293
598
    RETURN_IF_ERROR(_create_file_reader(false));
294
598
    RETURN_IF_ERROR(_create_decompressor());
295
598
    RETURN_IF_ERROR(_create_line_reader());
296
297
598
    _is_load = is_load;
298
598
    if (!_is_load) {
299
        // For query task, there are 2 slot mapping.
300
        // One is from file slot to values in line.
301
        //      eg, the file_slot_descs is k1, k3, k5, and values in line are k1, k2, k3, k4, k5
302
        //      the _col_idxs will save: 0, 2, 4
303
        // The other is from file slot to columns in output block
304
        //      eg, the file_slot_descs is k1, k3, k5, and columns in block are p1, k1, k3, k5
305
        //      where "p1" is the partition col which does not exist in file
306
        //      the _file_slot_idx_map will save: 1, 2, 3
307
0
        DCHECK(_params.__isset.column_idxs);
308
0
        _col_idxs = _params.column_idxs;
309
0
        int idx = 0;
310
0
        for (const auto& slot_info : _params.required_slots) {
311
0
            if (slot_info.is_file_slot) {
312
0
                _file_slot_idx_map.push_back(idx);
313
0
            }
314
0
            idx++;
315
0
        }
316
598
    } else {
317
        // For load task, the column order is same as file column order
318
598
        int i = 0;
319
598
        for (const auto& desc [[maybe_unused]] : _file_slot_descs) {
320
598
            _col_idxs.push_back(i++);
321
598
        }
322
598
    }
323
324
598
    _line_reader_eof = false;
325
598
    return Status::OK();
326
598
}
327
328
// ---- Unified init_reader(ReaderInitContext*) overrides ----
329
330
2.08k
Status CsvReader::_open_file_reader(ReaderInitContext* base_ctx) {
331
2.08k
    _start_offset = _range.start_offset;
332
2.08k
    if (_start_offset == 0) {
333
2.08k
        if (_params.__isset.file_attributes && _params.file_attributes.__isset.header_type &&
334
2.08k
            !_params.file_attributes.header_type.empty()) {
335
8
            std::string header_type = to_lower(_params.file_attributes.header_type);
336
8
            if (header_type == BeConsts::CSV_WITH_NAMES) {
337
5
                _skip_lines = 1;
338
5
            } else if (header_type == BeConsts::CSV_WITH_NAMES_AND_TYPES) {
339
3
                _skip_lines = 2;
340
3
            }
341
2.07k
        } else if (_params.file_attributes.__isset.skip_lines) {
342
2.07k
            _skip_lines = _params.file_attributes.skip_lines;
343
2.07k
        }
344
2.08k
    } else if (_start_offset != 0) {
345
0
        if ((_file_compress_type != TFileCompressType::PLAIN) ||
346
0
            (_file_compress_type == TFileCompressType::UNKNOWN &&
347
0
             _file_format_type != TFileFormatType::FORMAT_CSV_PLAIN)) {
348
0
            return Status::InternalError<false>("For now we do not support split compressed file");
349
0
        }
350
0
        int64_t pre_read_len = std::min(
351
0
                static_cast<int64_t>(_params.file_attributes.text_params.line_delimiter.size()),
352
0
                _start_offset);
353
0
        _start_offset -= pre_read_len;
354
0
        _size += pre_read_len;
355
0
        _skip_lines = 1;
356
0
    }
357
358
2.08k
    RETURN_IF_ERROR(_init_options());
359
2.08k
    RETURN_IF_ERROR(_create_file_reader(false));
360
2.08k
    return Status::OK();
361
2.08k
}
362
363
2.08k
Status CsvReader::_do_init_reader(ReaderInitContext* base_ctx) {
364
2.08k
    auto* ctx = checked_context_cast<CsvInitContext>(base_ctx);
365
2.08k
    _is_load = ctx->is_load;
366
367
2.08k
    _use_nullable_string_opt.resize(_file_slot_descs.size());
368
15.8k
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
369
13.7k
        auto data_type_ptr = _file_slot_descs[i]->get_data_type_ptr();
370
13.7k
        if (data_type_ptr->is_nullable() && is_string_type(data_type_ptr->get_primitive_type())) {
371
13.7k
            _use_nullable_string_opt[i] = 1;
372
13.7k
        }
373
13.7k
    }
374
375
2.08k
    RETURN_IF_ERROR(_create_decompressor());
376
2.08k
    RETURN_IF_ERROR(_create_line_reader());
377
378
2.08k
    if (!_is_load) {
379
2
        DCHECK(_params.__isset.column_idxs);
380
2
        _col_idxs = _params.column_idxs;
381
2
        int idx = 0;
382
6
        for (const auto& slot_info : _params.required_slots) {
383
6
            if (slot_info.is_file_slot) {
384
6
                _file_slot_idx_map.push_back(idx);
385
6
            }
386
6
            idx++;
387
6
        }
388
2.08k
    } else {
389
2.08k
        int i = 0;
390
13.7k
        for (const auto& desc [[maybe_unused]] : _file_slot_descs) {
391
13.7k
            _col_idxs.push_back(i++);
392
13.7k
        }
393
2.08k
    }
394
2.08k
    _line_reader_eof = false;
395
2.08k
    return Status::OK();
396
2.08k
}
397
398
7.14k
void CsvReader::set_batch_size(size_t batch_size) {
399
    // 0 means "not set" / "use default" for the row-based readers; we must
400
    // never let _batch_size be 0 because _do_get_next_block uses it as the
401
    // upper bound of a `while (rows < _batch_size)` loop and a 0 would make
402
    // the reader return empty blocks and incorrectly signal EOF.
403
7.14k
    _batch_size = std::max(batch_size, 1UL);
404
7.14k
}
405
406
// !FIXME: Here we should use MutableBlock
407
8.21k
Status CsvReader::_do_get_next_block(Block* block, size_t* read_rows, bool* eof) {
408
8.21k
    if (_line_reader_eof) {
409
2.30k
        *eof = true;
410
2.30k
        return Status::OK();
411
2.30k
    }
412
413
5.90k
    const size_t batch_size = _batch_size;
414
5.90k
    const auto max_block_bytes = _state->preferred_block_size_bytes();
415
5.90k
    size_t rows = 0;
416
417
5.90k
    bool success = false;
418
5.90k
    bool is_remove_bom = false;
419
5.90k
    if (_range.start_offset != 0 && _skip_lines > 0 && _enclose == 0 &&
420
5.90k
        _file_format_type == TFileFormatType::FORMAT_CSV_PLAIN) {
421
370
        auto* text_reader = assert_cast<NewPlainTextLineReader*>(_line_reader.get());
422
370
        RETURN_IF_ERROR(text_reader->skip_split_prefix(_range.start_offset, _line_delimiter,
423
370
                                                       &_line_reader_eof, _io_ctx));
424
370
        _skip_lines = 0;
425
370
        is_remove_bom = true;
426
370
    }
427
5.90k
    if (_push_down_agg_type == TPushAggOp::type::COUNT) {
428
1.04k
        while (rows < batch_size && !_line_reader_eof) {
429
646
            const uint8_t* ptr = nullptr;
430
646
            size_t size = 0;
431
646
            RETURN_IF_ERROR(_line_reader->read_line(&ptr, &size, &_line_reader_eof, _io_ctx));
432
433
            // _skip_lines == 0 means this line is the actual data beginning line for the entire file
434
            // is_remove_bom means _remove_bom should only execute once
435
646
            if (_skip_lines == 0 && !is_remove_bom) {
436
215
                ptr = _remove_bom(ptr, size);
437
215
                is_remove_bom = true;
438
215
            }
439
440
            // _skip_lines > 0 means we do not need to remove bom
441
646
            if (_skip_lines > 0) {
442
5
                _skip_lines--;
443
5
                is_remove_bom = true;
444
5
                continue;
445
5
            }
446
641
            if (size == 0) {
447
299
                if (!_line_reader_eof) {
448
0
                    if (_empty_line_as_record() || _state->is_read_csv_empty_line_as_null()) {
449
0
                        ++rows;
450
0
                    }
451
0
                }
452
                // Read empty line, continue
453
299
                continue;
454
299
            }
455
456
342
            RETURN_IF_ERROR(_validate_line(Slice(ptr, size), &success));
457
342
            ++rows;
458
342
        }
459
403
        auto mutable_columns_guard = block->mutate_columns_scoped();
460
403
        auto& mutate_columns = mutable_columns_guard.mutable_columns();
461
403
        for (auto& col : mutate_columns) {
462
403
            col->resize(rows);
463
403
        }
464
5.50k
    } else {
465
5.50k
        auto columns_guard = block->mutate_columns_scoped();
466
5.50k
        auto& columns = columns_guard.mutable_columns();
467
5.50k
        _reserve_nullable_string_columns(columns, batch_size);
468
10.8M
        while (rows < batch_size && !_line_reader_eof &&
469
10.8M
               (columns_byte_size(columns) < max_block_bytes)) {
470
10.7M
            const uint8_t* ptr = nullptr;
471
10.7M
            size_t size = 0;
472
10.7M
            RETURN_IF_ERROR(_line_reader->read_line(&ptr, &size, &_line_reader_eof, _io_ctx));
473
474
            // _skip_lines == 0 means this line is the actual data beginning line for the entire file
475
            // is_remove_bom means _remove_bom should only execute once
476
10.7M
            if (!is_remove_bom && _skip_lines == 0) {
477
5.28k
                ptr = _remove_bom(ptr, size);
478
5.28k
                is_remove_bom = true;
479
5.28k
            }
480
481
            // _skip_lines > 0 means we do not remove bom
482
10.7M
            if (_skip_lines > 0) {
483
57
                _skip_lines--;
484
57
                is_remove_bom = true;
485
57
                continue;
486
57
            }
487
10.7M
            if (size == 0) {
488
2.38k
                if (!_line_reader_eof) {
489
22
                    if (_empty_line_as_record()) {
490
0
                        Slice empty_line("", 0);
491
0
                        RETURN_IF_ERROR(_validate_line(empty_line, &success));
492
0
                        if (success) {
493
0
                            RETURN_IF_ERROR(_fill_dest_columns(empty_line, columns, &rows));
494
0
                        }
495
22
                    } else if (_state->is_read_csv_empty_line_as_null()) {
496
0
                        RETURN_IF_ERROR(_fill_empty_line(columns, &rows));
497
0
                    }
498
22
                }
499
                // Read empty line, continue
500
2.38k
                continue;
501
2.38k
            }
502
503
10.7M
            RETURN_IF_ERROR(_validate_line(Slice(ptr, size), &success));
504
10.7M
            if (!success) {
505
128
                continue;
506
128
            }
507
10.7M
            RETURN_IF_ERROR(_fill_dest_columns(Slice(ptr, size), columns, &rows));
508
10.7M
        }
509
5.50k
    }
510
511
5.89k
    *eof = (rows == 0);
512
5.89k
    *read_rows = rows;
513
514
5.89k
    return Status::OK();
515
5.90k
}
516
517
2.08k
Status CsvReader::_get_columns_impl(std::unordered_map<std::string, DataTypePtr>* name_to_type) {
518
13.7k
    for (const auto& slot : _file_slot_descs) {
519
13.7k
        name_to_type->emplace(slot->col_name(), slot->type());
520
13.7k
    }
521
2.08k
    return Status::OK();
522
2.08k
}
523
524
// init decompressor, file reader and line reader for parsing schema
525
403
Status CsvReader::init_schema_reader() {
526
403
    _start_offset = _range.start_offset;
527
403
    if (_start_offset != 0) {
528
0
        return Status::InvalidArgument(
529
0
                "start offset of TFileRangeDesc must be zero in get parsered schema");
530
0
    }
531
403
    if (_params.file_type == TFileType::FILE_BROKER) {
532
0
        return Status::InternalError<false>(
533
0
                "Getting parsered schema from csv file do not support stream load and broker "
534
0
                "load.");
535
0
    }
536
537
    // csv file without names line and types line.
538
403
    _read_line = 1;
539
403
    _is_parse_name = false;
540
541
403
    if (_params.__isset.file_attributes && _params.file_attributes.__isset.header_type &&
542
403
        !_params.file_attributes.header_type.empty()) {
543
48
        std::string header_type = to_lower(_params.file_attributes.header_type);
544
48
        if (header_type == BeConsts::CSV_WITH_NAMES) {
545
25
            _is_parse_name = true;
546
25
        } else if (header_type == BeConsts::CSV_WITH_NAMES_AND_TYPES) {
547
23
            _read_line = 2;
548
23
            _is_parse_name = true;
549
23
        }
550
48
    }
551
552
403
    RETURN_IF_ERROR(_init_options());
553
403
    RETURN_IF_ERROR(_create_file_reader(true));
554
403
    RETURN_IF_ERROR(_create_decompressor());
555
403
    RETURN_IF_ERROR(_create_line_reader());
556
403
    return Status::OK();
557
403
}
558
559
Status CsvReader::get_parsed_schema(std::vector<std::string>* col_names,
560
403
                                    std::vector<DataTypePtr>* col_types) {
561
403
    if (_read_line == 1) {
562
380
        if (!_is_parse_name) { //parse csv file without names and types
563
355
            size_t col_nums = 0;
564
355
            RETURN_IF_ERROR(_parse_col_nums(&col_nums));
565
3.00k
            for (size_t i = 0; i < col_nums; ++i) {
566
2.65k
                col_names->emplace_back("c" + std::to_string(i + 1));
567
2.65k
            }
568
354
        } else { // parse csv file with names
569
25
            RETURN_IF_ERROR(_parse_col_names(col_names));
570
25
        }
571
572
3.13k
        for (size_t j = 0; j < col_names->size(); ++j) {
573
2.75k
            col_types->emplace_back(
574
2.75k
                    DataTypeFactory::instance().create_data_type(PrimitiveType::TYPE_STRING, true));
575
2.75k
        }
576
379
    } else { // parse csv file with names and types
577
23
        RETURN_IF_ERROR(_parse_col_names(col_names));
578
23
        RETURN_IF_ERROR(_parse_col_types(col_names->size(), col_types));
579
23
    }
580
402
    return Status::OK();
581
403
}
582
583
173M
Status CsvReader::_deserialize_nullable_string(IColumn& column, Slice& slice) {
584
    // This is the per-row per-column hot path of CSV load (load reads every column as
585
    // nullable string). The column type was already verified by the checked assert_cast
586
    // in _reserve_nullable_string_columns at the beginning of the batch, so the casts
587
    // here can skip the release-build typeid check.
588
173M
    auto& null_column = assert_cast<ColumnNullable&, TypeCheckOnRelease::DISABLE>(column);
589
173M
    auto& string_column = assert_cast<ColumnString&, TypeCheckOnRelease::DISABLE>(
590
173M
            null_column.get_nested_column());
591
173M
    if (_empty_field_as_null && slice.size == 0) {
592
3
        string_column.insert_default();
593
3
        null_column.get_null_map_data().push_back(1);
594
3
        return Status::OK();
595
3
    }
596
173M
    if (_options.null_len > 0 && !(_options.converted_from_string && slice.trim_double_quotes())) {
597
170M
        if (slice.compare(Slice(_options.null_format, _options.null_len)) == 0) {
598
30.0k
            string_column.insert_default();
599
30.0k
            null_column.get_null_map_data().push_back(1);
600
30.0k
            return Status::OK();
601
30.0k
        }
602
170M
    }
603
    // Same as DataTypeStringSerDe::deserialize_one_cell_from_csv (which never fails),
604
    // written out here to skip the SerDe layer and its per-cell assert_cast.
605
173M
    if (_options.escape_char != 0) {
606
3.23k
        escape_string_for_csv(slice.data, &slice.size, _options.escape_char, _options.quote_char);
607
3.23k
    }
608
173M
    string_column.insert_data(slice.data, slice.size);
609
173M
    null_column.get_null_map_data().push_back(0);
610
173M
    return Status::OK();
611
173M
}
612
613
3.08k
Status CsvReader::_init_options() {
614
    // get column_separator and line_delimiter
615
3.08k
    _value_separator = _params.file_attributes.text_params.column_separator;
616
3.08k
    _value_separator_length = _value_separator.size();
617
3.08k
    _line_delimiter = _params.file_attributes.text_params.line_delimiter;
618
3.08k
    _line_delimiter_length = _line_delimiter.size();
619
3.08k
    if (_params.file_attributes.text_params.__isset.enclose) {
620
2.48k
        _enclose = _params.file_attributes.text_params.enclose;
621
2.48k
    }
622
3.08k
    if (_params.file_attributes.text_params.__isset.escape) {
623
2.48k
        _escape = _params.file_attributes.text_params.escape;
624
2.48k
    }
625
626
3.08k
    _trim_tailing_spaces =
627
3.08k
            (_state != nullptr && _state->trim_tailing_spaces_for_external_table_query());
628
629
3.08k
    _options.escape_char = _escape;
630
3.08k
    _options.quote_char = _enclose;
631
632
3.08k
    if (_params.file_attributes.text_params.collection_delimiter.empty()) {
633
3.08k
        _options.collection_delim = ',';
634
3.08k
    } else {
635
0
        _options.collection_delim = _params.file_attributes.text_params.collection_delimiter[0];
636
0
    }
637
3.08k
    if (_params.file_attributes.text_params.mapkv_delimiter.empty()) {
638
3.08k
        _options.map_key_delim = ':';
639
3.08k
    } else {
640
0
        _options.map_key_delim = _params.file_attributes.text_params.mapkv_delimiter[0];
641
0
    }
642
643
3.08k
    if (_params.file_attributes.text_params.__isset.null_format) {
644
0
        _options.null_format = _params.file_attributes.text_params.null_format.data();
645
0
        _options.null_len = _params.file_attributes.text_params.null_format.length();
646
0
    }
647
648
3.08k
    if (_params.file_attributes.__isset.trim_double_quotes) {
649
2.48k
        _trim_double_quotes = _params.file_attributes.trim_double_quotes;
650
2.48k
    }
651
3.08k
    _options.converted_from_string = _trim_double_quotes;
652
653
3.08k
    if (_state != nullptr) {
654
2.68k
        _keep_cr = _state->query_options().keep_carriage_return;
655
2.68k
    }
656
657
3.08k
    if (_params.file_attributes.text_params.__isset.empty_field_as_null) {
658
2.48k
        _empty_field_as_null = _params.file_attributes.text_params.empty_field_as_null;
659
2.48k
    }
660
3.08k
    return Status::OK();
661
3.08k
}
662
663
3.08k
Status CsvReader::_create_decompressor() {
664
3.08k
    if (_file_compress_type != TFileCompressType::UNKNOWN) {
665
3.08k
        RETURN_IF_ERROR(Decompressor::create_decompressor(_file_compress_type, &_decompressor));
666
3.08k
    } else {
667
1
        RETURN_IF_ERROR(Decompressor::create_decompressor(_file_format_type, &_decompressor));
668
1
    }
669
670
3.08k
    return Status::OK();
671
3.08k
}
672
673
3.08k
Status CsvReader::_create_file_reader(bool need_schema) {
674
3.08k
    if (_params.file_type == TFileType::FILE_STREAM) {
675
2.12k
        RETURN_IF_ERROR(FileFactory::create_pipe_reader(_range.load_id, &_file_reader, _state,
676
2.12k
                                                        need_schema));
677
2.12k
    } else {
678
958
        _file_description.mtime = _range.__isset.modification_time ? _range.modification_time : 0;
679
958
        io::FileReaderOptions reader_options = FileFactory::get_reader_options(
680
958
                _state ? _state->query_options() : _default_query_options, _file_description);
681
958
        io::FileReaderSPtr file_reader;
682
958
        if (_io_ctx_holder) {
683
360
            file_reader = DORIS_TRY(io::DelegateReader::create_file_reader(
684
360
                    _profile, _system_properties, _file_description, reader_options,
685
360
                    io::DelegateReader::AccessMode::SEQUENTIAL,
686
360
                    std::static_pointer_cast<const io::IOContext>(_io_ctx_holder),
687
360
                    io::PrefetchRange(_range.start_offset, _range.start_offset + _range.size)));
688
598
        } else {
689
598
            file_reader = DORIS_TRY(io::DelegateReader::create_file_reader(
690
598
                    _profile, _system_properties, _file_description, reader_options,
691
598
                    io::DelegateReader::AccessMode::SEQUENTIAL, _io_ctx,
692
598
                    io::PrefetchRange(_range.start_offset, _range.start_offset + _range.size)));
693
598
        }
694
958
        _file_reader = _io_ctx && _io_ctx->file_reader_stats
695
958
                               ? std::make_shared<io::TracingFileReader>(std::move(file_reader),
696
360
                                                                         _io_ctx->file_reader_stats)
697
958
                               : file_reader;
698
958
    }
699
3.08k
    if (_file_reader->size() == 0 && _params.file_type != TFileType::FILE_STREAM &&
700
3.08k
        _params.file_type != TFileType::FILE_BROKER) {
701
0
        return Status::EndOfFile("init reader failed, empty csv file: " + _range.path);
702
0
    }
703
3.08k
    return Status::OK();
704
3.08k
}
705
706
3.08k
Status CsvReader::_create_line_reader() {
707
3.08k
    std::shared_ptr<TextLineReaderContextIf> text_line_reader_ctx;
708
3.08k
    if (_enclose == 0) {
709
2.99k
        text_line_reader_ctx = std::make_shared<PlainTextLineReaderCtx>(
710
2.99k
                _line_delimiter, _line_delimiter_length, _keep_cr);
711
2.99k
        _fields_splitter = std::make_unique<PlainCsvTextFieldSplitter>(
712
2.99k
                _trim_tailing_spaces, false, _value_separator, _value_separator_length, -1);
713
714
2.99k
    } else {
715
        // in load task, the _file_slot_descs is empty vector, so we need to set col_sep_num to 0
716
92
        size_t col_sep_num = _file_slot_descs.size() > 1 ? _file_slot_descs.size() - 1 : 0;
717
92
        _enclose_reader_ctx = std::make_shared<EncloseCsvLineReaderCtx>(
718
92
                _line_delimiter, _line_delimiter_length, _value_separator, _value_separator_length,
719
92
                col_sep_num, _enclose, _escape, _keep_cr);
720
92
        text_line_reader_ctx = _enclose_reader_ctx;
721
722
92
        _fields_splitter = std::make_unique<EncloseCsvTextFieldSplitter>(
723
92
                _trim_tailing_spaces, true, _enclose_reader_ctx, _value_separator_length, _enclose);
724
92
    }
725
3.08k
    switch (_file_format_type) {
726
3.00k
    case TFileFormatType::FORMAT_CSV_PLAIN:
727
3.00k
        [[fallthrough]];
728
3.00k
    case TFileFormatType::FORMAT_CSV_GZ:
729
3.00k
        [[fallthrough]];
730
3.00k
    case TFileFormatType::FORMAT_CSV_BZ2:
731
3.00k
        [[fallthrough]];
732
3.00k
    case TFileFormatType::FORMAT_CSV_LZ4FRAME:
733
3.00k
        [[fallthrough]];
734
3.00k
    case TFileFormatType::FORMAT_CSV_LZ4BLOCK:
735
3.00k
        [[fallthrough]];
736
3.00k
    case TFileFormatType::FORMAT_CSV_LZOP:
737
3.00k
        [[fallthrough]];
738
3.00k
    case TFileFormatType::FORMAT_CSV_SNAPPYBLOCK:
739
3.00k
        [[fallthrough]];
740
3.00k
    case TFileFormatType::FORMAT_CSV_DEFLATE:
741
3.00k
        _line_reader =
742
3.00k
                NewPlainTextLineReader::create_unique(_profile, _file_reader, _decompressor.get(),
743
3.00k
                                                      text_line_reader_ctx, _size, _start_offset);
744
745
3.00k
        break;
746
78
    case TFileFormatType::FORMAT_PROTO:
747
78
        _fields_splitter = std::make_unique<CsvProtoFieldSplitter>();
748
78
        _line_reader = NewPlainBinaryLineReader::create_unique(_file_reader);
749
78
        break;
750
0
    default:
751
0
        return Status::InternalError<false>(
752
0
                "Unknown format type, cannot init line reader in csv reader, type={}",
753
0
                _file_format_type);
754
3.08k
    }
755
3.08k
    return Status::OK();
756
3.08k
}
757
758
0
Status CsvReader::_deserialize_one_cell(DataTypeSerDeSPtr serde, IColumn& column, Slice& slice) {
759
0
    return serde->deserialize_one_cell_from_csv(column, slice, _options);
760
0
}
761
762
Status CsvReader::_fill_dest_columns(const Slice& line, std::vector<MutableColumnPtr>& columns,
763
10.8M
                                     size_t* rows) {
764
10.8M
    bool is_success = false;
765
766
10.8M
    RETURN_IF_ERROR(_line_split_to_values(line, &is_success));
767
10.8M
    if (UNLIKELY(!is_success)) {
768
        // If not success, which means we met an invalid row, filter this row and return.
769
539
        return Status::OK();
770
539
    }
771
772
175M
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
773
164M
        int col_idx = _col_idxs[i];
774
        // col idx is out of range, fill with null format
775
164M
        auto value = col_idx < _split_values.size()
776
165M
                             ? _split_values[col_idx]
777
18.4E
                             : Slice(_options.null_format, _options.null_len);
778
779
164M
        IColumn* col_ptr = columns[i].get();
780
164M
        if (!_is_load) {
781
0
            col_ptr = columns[_file_slot_idx_map[i]].get();
782
0
        }
783
784
175M
        if (_use_nullable_string_opt[i]) {
785
            // For load task, we always read "string" from file.
786
            // So serdes[i] here must be DataTypeNullableSerDe, and DataTypeNullableSerDe -> nested_serde must be DataTypeStringSerDe.
787
            // So we use deserialize_nullable_string and stringSerDe to reduce virtual function calls.
788
175M
            RETURN_IF_ERROR(_deserialize_nullable_string(*col_ptr, value));
789
18.4E
        } else {
790
18.4E
            RETURN_IF_ERROR(_deserialize_one_cell(_serdes[i], *col_ptr, value));
791
18.4E
        }
792
164M
    }
793
10.8M
    ++(*rows);
794
795
10.8M
    return Status::OK();
796
10.8M
}
797
798
void CsvReader::_reserve_nullable_string_columns(std::vector<MutableColumnPtr>& columns,
799
5.50k
                                                 size_t batch_size) {
800
67.6k
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
801
62.1k
        if (!_use_nullable_string_opt[i]) {
802
0
            continue;
803
0
        }
804
18.4E
        IColumn* col_ptr = _is_load ? columns[i].get() : columns[_file_slot_idx_map[i]].get();
805
        // The checked casts here (once per batch) guarantee the column types for the
806
        // unchecked per-row casts in _deserialize_nullable_string.
807
62.1k
        auto& null_column = assert_cast<ColumnNullable&>(*col_ptr);
808
62.1k
        auto& string_column = assert_cast<ColumnString&>(null_column.get_nested_column());
809
        // Reserve up front so the per-row loop does not pay for incremental growth.
810
        // The string chars are not reserved because their total size is unpredictable.
811
62.1k
        string_column.get_offsets().reserve(string_column.size() + batch_size);
812
62.1k
        null_column.get_null_map_data().reserve(null_column.get_null_map_data().size() +
813
62.1k
                                                batch_size);
814
62.1k
    }
815
5.50k
}
816
817
0
Status CsvReader::_fill_empty_line(std::vector<MutableColumnPtr>& columns, size_t* rows) {
818
0
    for (int i = 0; i < _file_slot_descs.size(); ++i) {
819
0
        IColumn* col_ptr = columns[i].get();
820
0
        if (!_is_load) {
821
0
            col_ptr = columns[_file_slot_idx_map[i]].get();
822
0
        }
823
0
        auto& null_column = assert_cast<ColumnNullable&>(*col_ptr);
824
0
        null_column.insert_data(nullptr, 0);
825
0
    }
826
0
    ++(*rows);
827
0
    return Status::OK();
828
0
}
829
830
10.8M
Status CsvReader::_validate_line(const Slice& line, bool* success) {
831
10.8M
    if (!_is_proto_format && !validate_utf8(_params, line.data, line.size)) {
832
128
        if (!_is_load) {
833
0
            return Status::InternalError<false>("Only support csv data in utf8 codec");
834
128
        } else {
835
128
            _counter->num_rows_filtered++;
836
128
            *success = false;
837
128
            RETURN_IF_ERROR(_state->append_error_msg_to_file(
838
128
                    [&]() -> std::string { return std::string(line.data, line.size); },
839
128
                    [&]() -> std::string {
840
128
                        return "Invalid file encoding: all CSV files must be UTF-8 encoded";
841
128
                    }));
842
128
            return Status::OK();
843
128
        }
844
128
    }
845
10.8M
    *success = true;
846
10.8M
    return Status::OK();
847
10.8M
}
848
849
10.8M
Status CsvReader::_line_split_to_values(const Slice& line, bool* success) {
850
10.8M
    _split_line(line);
851
852
10.8M
    if (_is_load) {
853
        // Only check for load task. For query task, the non exist column will be filled "null".
854
        // if actual column number in csv file is not equal to _file_slot_descs.size()
855
        // then filter this line.
856
10.8M
        bool ignore_col = false;
857
10.8M
        ignore_col = _params.__isset.file_attributes &&
858
10.8M
                     _params.file_attributes.__isset.ignore_csv_redundant_col &&
859
10.8M
                     _params.file_attributes.ignore_csv_redundant_col;
860
861
10.8M
        if ((!ignore_col && _split_values.size() != _file_slot_descs.size()) ||
862
10.8M
            (ignore_col && _split_values.size() < _file_slot_descs.size())) {
863
544
            _counter->num_rows_filtered++;
864
544
            *success = false;
865
544
            RETURN_IF_ERROR(_state->append_error_msg_to_file(
866
544
                    [&]() -> std::string { return std::string(line.data, line.size); },
867
544
                    [&]() -> std::string {
868
544
                        fmt::memory_buffer error_msg;
869
544
                        fmt::format_to(error_msg,
870
544
                                       "Column count mismatch: expected {}, but found {}",
871
544
                                       _file_slot_descs.size(), _split_values.size());
872
544
                        std::string escaped_separator =
873
544
                                std::regex_replace(_value_separator, std::regex("\t"), "\\t");
874
544
                        std::string escaped_delimiter =
875
544
                                std::regex_replace(_line_delimiter, std::regex("\n"), "\\n");
876
544
                        fmt::format_to(error_msg, " (sep:{} delim:{}", escaped_separator,
877
544
                                       escaped_delimiter);
878
544
                        if (_enclose != 0) {
879
544
                            fmt::format_to(error_msg, " encl:{}", _enclose);
880
544
                        }
881
544
                        if (_escape != 0) {
882
544
                            fmt::format_to(error_msg, " esc:{}", _escape);
883
544
                        }
884
544
                        fmt::format_to(error_msg, ")");
885
544
                        return fmt::to_string(error_msg);
886
544
                    }));
887
539
            return Status::OK();
888
544
        }
889
10.8M
    }
890
891
10.8M
    *success = true;
892
10.8M
    return Status::OK();
893
10.8M
}
894
895
10.8M
void CsvReader::_split_line(const Slice& line) {
896
10.8M
    _split_values.clear();
897
10.8M
    _fields_splitter->split_line(line, &_split_values);
898
10.8M
}
899
900
355
Status CsvReader::_parse_col_nums(size_t* col_nums) {
901
355
    const uint8_t* ptr = nullptr;
902
355
    size_t size = 0;
903
355
    RETURN_IF_ERROR(_line_reader->read_line(&ptr, &size, &_line_reader_eof, _io_ctx));
904
355
    if (size == 0) {
905
0
        return Status::InternalError<false>(
906
0
                "The first line is empty, can not parse column numbers");
907
0
    }
908
355
    if (!validate_utf8(_params, reinterpret_cast<const char*>(ptr), size)) {
909
1
        return Status::InternalError<false>("Only support csv data in utf8 codec");
910
1
    }
911
354
    ptr = _remove_bom(ptr, size);
912
354
    _split_line(Slice(ptr, size));
913
354
    *col_nums = _split_values.size();
914
354
    return Status::OK();
915
355
}
916
917
48
Status CsvReader::_parse_col_names(std::vector<std::string>* col_names) {
918
48
    const uint8_t* ptr = nullptr;
919
48
    size_t size = 0;
920
    // no use of _line_reader_eof
921
48
    RETURN_IF_ERROR(_line_reader->read_line(&ptr, &size, &_line_reader_eof, _io_ctx));
922
48
    if (size == 0) {
923
0
        return Status::InternalError<false>("The first line is empty, can not parse column names");
924
0
    }
925
48
    if (!validate_utf8(_params, reinterpret_cast<const char*>(ptr), size)) {
926
0
        return Status::InternalError<false>("Only support csv data in utf8 codec");
927
0
    }
928
48
    ptr = _remove_bom(ptr, size);
929
48
    _split_line(Slice(ptr, size));
930
187
    for (auto _split_value : _split_values) {
931
187
        col_names->emplace_back(_split_value.to_string());
932
187
    }
933
48
    return Status::OK();
934
48
}
935
936
// TODO(ftw): parse type
937
23
Status CsvReader::_parse_col_types(size_t col_nums, std::vector<DataTypePtr>* col_types) {
938
    // delete after.
939
113
    for (size_t i = 0; i < col_nums; ++i) {
940
90
        col_types->emplace_back(make_nullable(std::make_shared<DataTypeString>()));
941
90
    }
942
943
    // 1. check _line_reader_eof
944
    // 2. read line
945
    // 3. check utf8
946
    // 4. check size
947
    // 5. check _split_values.size must equal to col_nums.
948
    // 6. fill col_types
949
23
    return Status::OK();
950
23
}
951
952
5.90k
const uint8_t* CsvReader::_remove_bom(const uint8_t* ptr, size_t& size) {
953
5.90k
    if (size >= 3 && ptr[0] == 0xEF && ptr[1] == 0xBB && ptr[2] == 0xBF) {
954
5
        LOG(INFO) << "remove bom";
955
5
        constexpr size_t bom_size = 3;
956
5
        size -= bom_size;
957
        // In enclose mode, column_sep_positions were computed on the original line
958
        // (including BOM). After shifting the pointer, we must adjust those positions
959
        // so they remain correct relative to the new start.
960
5
        if (_enclose_reader_ctx) {
961
1
            _enclose_reader_ctx->adjust_column_sep_positions(bom_size);
962
1
        }
963
5
        return ptr + bom_size;
964
5
    }
965
5.89k
    return ptr;
966
5.90k
}
967
968
2.08k
Status CsvReader::close() {
969
2.08k
    if (_line_reader) {
970
2.08k
        _line_reader->close();
971
2.08k
    }
972
973
2.08k
    if (_file_reader) {
974
2.08k
        RETURN_IF_ERROR(_file_reader->close());
975
2.08k
    }
976
977
2.08k
    return Status::OK();
978
2.08k
}
979
980
} // namespace doris